Skip to content

Commit a2e3300

Browse files
committed
start
1 parent b8d4811 commit a2e3300

File tree

16 files changed

+36
-26
lines changed

16 files changed

+36
-26
lines changed

core/src/main/java/com/dtstack/flink/sql/source/AbsDeserialization.java renamed to core/src/main/java/com/dtstack/flink/sql/format/AbsDeserialization.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,9 +16,10 @@
1616
* limitations under the License.
1717
*/
1818

19-
package com.dtstack.flink.sql.source;
19+
package com.dtstack.flink.sql.format;
2020

2121
import com.dtstack.flink.sql.metric.MetricConstant;
22+
import com.dtstack.flink.sql.source.JsonDataParser;
2223
import org.apache.flink.api.common.functions.RuntimeContext;
2324
import org.apache.flink.api.common.serialization.AbstractDeserializationSchema;
2425
import org.apache.flink.metrics.Counter;

kafka/kafka-source/src/main/java/com/dtstack/flink/sql/source/kafka/CustomerJsonDeserialization.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@
1919
package com.dtstack.flink.sql.source.kafka;
2020

2121

22-
import com.dtstack.flink.sql.source.AbsDeserialization;
22+
import com.dtstack.flink.sql.format.AbsDeserialization;
2323
import com.dtstack.flink.sql.source.JsonDataParser;
2424
import com.dtstack.flink.sql.source.kafka.metric.KafkaTopicPartitionLagMetric;
2525
import com.dtstack.flink.sql.table.TableInfo;

kafka/kafka-source/src/main/java/com/dtstack/flink/sql/source/kafka/CustomerKafkaConsumer.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818

1919
package com.dtstack.flink.sql.source.kafka;
2020

21-
import com.dtstack.flink.sql.source.AbsDeserialization;
21+
import com.dtstack.flink.sql.format.AbsDeserialization;
2222
import org.apache.flink.metrics.MetricGroup;
2323
import org.apache.flink.streaming.api.functions.AssignerWithPeriodicWatermarks;
2424
import org.apache.flink.streaming.api.functions.AssignerWithPunctuatedWatermarks;

kafka08/kafka08-source/src/main/java/com/dtstack/flink/sql/source/kafka/consumer/CustomerCsvConsumer.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818

1919
package com.dtstack.flink.sql.source.kafka.consumer;
2020

21-
import com.dtstack.flink.sql.source.AbsDeserialization;
21+
import com.dtstack.flink.sql.format.AbsDeserialization;
2222
import com.dtstack.flink.sql.source.kafka.deserialization.CustomerCsvDeserialization;
2323
import org.apache.flink.streaming.api.functions.source.SourceFunction;
2424
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer08;

kafka08/kafka08-source/src/main/java/com/dtstack/flink/sql/source/kafka/consumer/CustomerJsonConsumer.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818

1919
package com.dtstack.flink.sql.source.kafka.consumer;
2020

21-
import com.dtstack.flink.sql.source.AbsDeserialization;
21+
import com.dtstack.flink.sql.format.AbsDeserialization;
2222
import com.dtstack.flink.sql.source.kafka.deserialization.CustomerJsonDeserialization;
2323
import org.apache.flink.streaming.api.functions.source.SourceFunction;
2424
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer08;

kafka08/kafka08-source/src/main/java/com/dtstack/flink/sql/source/kafka/deserialization/CustomerCommonDeserialization.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818

1919
package com.dtstack.flink.sql.source.kafka.deserialization;
2020

21-
import com.dtstack.flink.sql.source.AbsDeserialization;
21+
import com.dtstack.flink.sql.format.AbsDeserialization;
2222
import org.apache.flink.api.common.typeinfo.TypeInformation;
2323
import org.apache.flink.api.common.typeinfo.Types;
2424
import org.apache.flink.api.java.typeutils.RowTypeInfo;

kafka08/kafka08-source/src/main/java/com/dtstack/flink/sql/source/kafka/deserialization/CustomerCsvDeserialization.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
package com.dtstack.flink.sql.source.kafka.deserialization;
2222

2323

24-
import com.dtstack.flink.sql.source.AbsDeserialization;
24+
import com.dtstack.flink.sql.format.AbsDeserialization;
2525
import com.dtstack.flink.sql.util.DtStringUtil;
2626
import org.apache.flink.api.common.typeinfo.TypeInformation;
2727
import org.apache.flink.api.java.typeutils.RowTypeInfo;

kafka08/kafka08-source/src/main/java/com/dtstack/flink/sql/source/kafka/deserialization/CustomerJsonDeserialization.java

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
package com.dtstack.flink.sql.source.kafka.deserialization;
2222

2323

24-
import com.dtstack.flink.sql.source.AbsDeserialization;
24+
import com.dtstack.flink.sql.format.AbsDeserialization;
2525
import org.apache.flink.api.common.typeinfo.TypeInformation;
2626
import org.apache.flink.api.java.typeutils.RowTypeInfo;
2727
import com.fasterxml.jackson.databind.JsonNode;
@@ -33,11 +33,7 @@
3333
import org.slf4j.LoggerFactory;
3434

3535
import java.io.IOException;
36-
import java.lang.reflect.Field;
3736
import java.util.Iterator;
38-
import java.util.Set;
39-
40-
import static com.dtstack.flink.sql.metric.MetricConstant.*;
4137

4238
/**
4339
* json string parsing custom

kafka09/kafka09-source/src/main/java/com/dtstack/flink/sql/source/kafka/CustomerJsonDeserialization.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
package com.dtstack.flink.sql.source.kafka;
2222

2323

24-
import com.dtstack.flink.sql.source.AbsDeserialization;
24+
import com.dtstack.flink.sql.format.AbsDeserialization;
2525
import com.dtstack.flink.sql.source.JsonDataParser;
2626
import com.dtstack.flink.sql.source.kafka.metric.KafkaTopicPartitionLagMetric;
2727
import com.dtstack.flink.sql.table.TableInfo;

kafka09/kafka09-source/src/main/java/com/dtstack/flink/sql/source/kafka/CustomerKafka09Consumer.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818

1919
package com.dtstack.flink.sql.source.kafka;
2020

21-
import com.dtstack.flink.sql.source.AbsDeserialization;
21+
import com.dtstack.flink.sql.format.AbsDeserialization;
2222
import org.apache.flink.metrics.MetricGroup;
2323
import org.apache.flink.streaming.api.functions.AssignerWithPeriodicWatermarks;
2424
import org.apache.flink.streaming.api.functions.AssignerWithPunctuatedWatermarks;

0 commit comments

Comments
 (0)