diff --git a/maven/flink-quickstart.iml b/maven/flink-quickstart.iml
index e9231310756e18deb654f127aaa9eca9668dbfb7..e0e5543b6613aa451b916f44a9aad7311dbe1b75 100644
--- a/maven/flink-quickstart.iml
+++ b/maven/flink-quickstart.iml
@@ -52,6 +52,17 @@
+
+
+
+
+
+
+
+
+
+
+
diff --git a/maven/pom.xml b/maven/pom.xml
index 7ee292bdf0e34006ffd2514f076bf5149360fcaf..93a52ad9297dba2d7c3952cd2decd0d539975ef5 100644
--- a/maven/pom.xml
+++ b/maven/pom.xml
@@ -66,6 +66,18 @@ under the License.
${flink.version}
provided
+
+ org.apache.flink
+ flink-table-planner_2.11
+ 1.9.1
+
+
+
+ org.apache.flink
+ flink-table-api-java-bridge_2.11
+ 1.9.1
+
+
com.google.guava
guava
diff --git a/maven/src/main/java/com/duanledexianxianxian/maven/flink/common/constant/AppConstants.java b/maven/src/main/java/com/duanledexianxianxian/maven/flink/common/constant/AppConstants.java
index 83a965a0606ca9b6d420ccce9a6c1a0eab28749e..fa21e3ed202776646a49edc1e85393e32c1e6f2b 100644
--- a/maven/src/main/java/com/duanledexianxianxian/maven/flink/common/constant/AppConstants.java
+++ b/maven/src/main/java/com/duanledexianxianxian/maven/flink/common/constant/AppConstants.java
@@ -7,7 +7,7 @@ public class AppConstants {
/**
* 配置文件
*/
- public final static String PROPERTIES_FILE_NAME = "application.properties";
+ public final static String PROPERTIES_FILE_NAME = "/application.properties";
public static final String STREAM_PARALLELISM = "stream.parallelism";
diff --git a/maven/src/main/java/com/duanledexianxianxian/maven/flink/common/schemas/SimplekafkaStringSchema.java b/maven/src/main/java/com/duanledexianxianxian/maven/flink/common/schemas/SimplekafkaStringSchema.java
new file mode 100644
index 0000000000000000000000000000000000000000..d5e54f62e4e283a67d2bee83d4c1f9f58a5fed07
--- /dev/null
+++ b/maven/src/main/java/com/duanledexianxianxian/maven/flink/common/schemas/SimplekafkaStringSchema.java
@@ -0,0 +1,45 @@
+package com.duanledexianxianxian.maven.flink.common.schemas;
+
+import com.alibaba.fastjson.JSONObject;
+import org.apache.flink.api.common.typeinfo.TypeInformation;
+import org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema;
+import org.apache.flink.streaming.connectors.kafka.KafkaSerializationSchema;
+import org.apache.kafka.clients.consumer.ConsumerRecord;
+import org.apache.kafka.clients.producer.ProducerRecord;
+
+import javax.annotation.Nullable;
+
+/**
+ * @author Administrator
+ */
+public class SimplekafkaStringSchema implements KafkaSerializationSchema