This is an automated email from the ASF dual-hosted git repository.

dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 2f56b3244 [INLONG-7666][Sort] Add a flag to determine whether the data 
is incremental data in all migrate (#7669)
2f56b3244 is described below

commit 2f56b32440f6c173d0a66506474df9727b5f4bac
Author: Schnapps <[email protected]>
AuthorDate: Mon Mar 27 12:55:12 2023 +0800

    [INLONG-7666][Sort] Add a flag to determine whether the data is incremental 
data in all migrate (#7669)
---
 .../inlong/sort/cdc/mysql/source/MySqlSource.java  |   3 +-
 .../sort/cdc/mysql/source/MySqlSourceBuilder.java  |   5 +
 .../cdc/mysql/source/config/MySqlSourceConfig.java |   9 +-
 .../source/config/MySqlSourceConfigFactory.java    |   9 +-
 .../mysql/source/config/MySqlSourceOptions.java    |   7 ++
 .../mysql/source/reader/MySqlRecordEmitter.java    |  25 +++--
 .../sort/cdc/mysql/source/utils/RecordUtils.java   |  13 +++
 .../cdc/mysql/table/MySqlReadableMetadata.java     |   7 +-
 .../mysql/table/MySqlTableInlongSourceFactory.java |   7 +-
 .../sort/cdc/mysql/table/MySqlTableSource.java     |  16 ++-
 .../apache/inlong/sort/parser/AllMigrateTest.java  |   1 +
 .../inlong/sort/formats/json/canal/CanalJson.java  | 124 ++++-----------------
 .../sort/formats/json/debezium/DebeziumJson.java   |   1 +
 13 files changed, 102 insertions(+), 125 deletions(-)

diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/MySqlSource.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/MySqlSource.java
index ea79e14b8..47fd75d4e 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/MySqlSource.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/MySqlSource.java
@@ -164,7 +164,8 @@ public class MySqlSource<T>
                 new MySqlRecordEmitter<>(
                         deserializationSchema,
                         sourceReaderMetrics,
-                        sourceConfig.isIncludeSchemaChanges()),
+                        sourceConfig.isIncludeSchemaChanges(),
+                        sourceConfig.isIncludeIncremental()),
                 readerContext.getConfiguration(),
                 mySqlSourceReaderContext,
                 sourceConfig, sourceReaderMetrics);
diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/MySqlSourceBuilder.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/MySqlSourceBuilder.java
index 86b9991ee..8cfaa0bae 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/MySqlSourceBuilder.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/MySqlSourceBuilder.java
@@ -270,4 +270,9 @@ public class MySqlSourceBuilder<T> {
     public MySqlSource<T> build() {
         return new MySqlSource<>(configFactory, checkNotNull(deserializer));
     }
+
+    public MySqlSourceBuilder<T> includeIncremental(boolean 
includeIncremental) {
+        this.configFactory.includeIncremental(includeIncremental);
+        return this;
+    }
 }
diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceConfig.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceConfig.java
index 4418c8576..926bbebce 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceConfig.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceConfig.java
@@ -69,6 +69,7 @@ public class MySqlSourceConfig implements Serializable {
 
     private final String inlongMetric;
     private final String inlongAudit;
+    private final boolean includeIncremental;
 
     MySqlSourceConfig(
             String hostname,
@@ -93,7 +94,8 @@ public class MySqlSourceConfig implements Serializable {
             Properties dbzProperties,
             Properties jdbcProperties,
             String inlongMetric,
-            String inlongAudit) {
+            String inlongAudit,
+            boolean includeIncremental) {
         this.hostname = checkNotNull(hostname);
         this.port = port;
         this.username = checkNotNull(username);
@@ -119,6 +121,7 @@ public class MySqlSourceConfig implements Serializable {
         this.jdbcProperties = jdbcProperties;
         this.inlongMetric = inlongMetric;
         this.inlongAudit = inlongAudit;
+        this.includeIncremental = includeIncremental;
     }
 
     public String getHostname() {
@@ -225,4 +228,8 @@ public class MySqlSourceConfig implements Serializable {
     public String getInlongAudit() {
         return inlongAudit;
     }
+
+    public boolean isIncludeIncremental() {
+        return includeIncremental;
+    }
 }
diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceConfigFactory.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceConfigFactory.java
index 10fe5f636..82b90d561 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceConfigFactory.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceConfigFactory.java
@@ -75,6 +75,7 @@ public class MySqlSourceConfigFactory implements Serializable 
{
 
     private String inlongMetric;
     private String inlongAudit;
+    private boolean includeIncremental;
 
     public MySqlSourceConfigFactory inlongMetric(String inlongMetric) {
         this.inlongMetric = inlongMetric;
@@ -86,6 +87,11 @@ public class MySqlSourceConfigFactory implements 
Serializable {
         return this;
     }
 
+    public MySqlSourceConfigFactory includeIncremental(boolean 
includeIncremental) {
+        this.includeIncremental = includeIncremental;
+        return this;
+    }
+
     public MySqlSourceConfigFactory hostname(String hostname) {
         this.hostname = hostname;
         return this;
@@ -369,6 +375,7 @@ public class MySqlSourceConfigFactory implements 
Serializable {
                 props,
                 jdbcProperties,
                 inlongMetric,
-                inlongAudit);
+                inlongAudit,
+                includeIncremental);
     }
 }
diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceOptions.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceOptions.java
index f770fac71..0a09dd400 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceOptions.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/config/MySqlSourceOptions.java
@@ -206,6 +206,13 @@ public class MySqlSourceOptions {
                     .defaultValue(false)
                     .withDescription("Whether migrate all databases");
 
+    public static final ConfigOption<Boolean> INCLUDE_INCREMENTAL =
+            ConfigOptions.key("include-incremental")
+                    .booleanType()
+                    .defaultValue(false)
+                    .withDescription("Whether include a incremental flag in 
data "
+                            + "when migrating all databases");
+
     // 
----------------------------------------------------------------------------
     // experimental options, won't add them to documentation
     // 
----------------------------------------------------------------------------
diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/reader/MySqlRecordEmitter.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/reader/MySqlRecordEmitter.java
index e9427ec3c..e36b0d80f 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/reader/MySqlRecordEmitter.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/reader/MySqlRecordEmitter.java
@@ -52,6 +52,7 @@ import static 
org.apache.inlong.sort.cdc.mysql.source.utils.RecordUtils.isHeartb
 import static 
org.apache.inlong.sort.cdc.mysql.source.utils.RecordUtils.isHighWatermarkEvent;
 import static 
org.apache.inlong.sort.cdc.mysql.source.utils.RecordUtils.isSchemaChangeEvent;
 import static 
org.apache.inlong.sort.cdc.mysql.source.utils.RecordUtils.isWatermarkEvent;
+import static 
org.apache.inlong.sort.cdc.mysql.source.utils.RecordUtils.toSnapshotRecord;
 
 /**
  * The {@link RecordEmitter} implementation for {@link MySqlSourceReader}.
@@ -80,25 +81,30 @@ public final class MySqlRecordEmitter<T>
     private volatile long snapEarliestTime = 0L;
     private volatile long snapProcessTime = 0L;
 
+    private boolean includeIncremental;
+
     public MySqlRecordEmitter(
             DebeziumDeserializationSchema<T> debeziumDeserializationSchema,
             MySqlSourceReaderMetrics sourceReaderMetrics,
-            boolean includeSchemaChanges) {
+            boolean includeSchemaChanges, boolean includeIncremental) {
         this.debeziumDeserializationSchema = debeziumDeserializationSchema;
         this.sourceReaderMetrics = sourceReaderMetrics;
         this.includeSchemaChanges = includeSchemaChanges;
         this.outputCollector = new OutputCollector<>();
+        this.includeIncremental = includeIncremental;
     }
 
     @Override
     public void emitRecord(SourceRecord element, SourceOutput<T> output, 
MySqlSplitState splitState)
             throws Exception {
+
         if (isWatermarkEvent(element)) {
             BinlogOffset watermark = getWatermark(element);
             if (isHighWatermarkEvent(element) && 
splitState.isSnapshotSplitState()) {
                 splitState.asSnapshotSplitState().setHighWatermark(watermark);
             }
         } else if (isSchemaChangeEvent(element) && 
splitState.isBinlogSplitState()) {
+            updateSnapshotRecord(element, splitState);
             HistoryRecord historyRecord = getHistoryRecord(element);
             Array tableChanges =
                     
historyRecord.document().getArray(HistoryRecord.Fields.TABLE_CHANGES);
@@ -112,15 +118,6 @@ public final class MySqlRecordEmitter<T>
                 emitElement(element, output, null);
             }
         } else if (isDataChangeRecord(element)) {
-            // updateStartingOffsetForSplit(splitState, element);
-            // reportMetrics(element);
-            //
-            // final Map<TableId, TableChange> tableSchemas =
-            // splitState.getMySQLSplit().getTableSchemas();
-            // final TableChange tableSchema =
-            // tableSchemas.getOrDefault(getTableId(element), null);
-            //
-            // emitElement(element, output, tableSchema);
             if (splitState.isBinlogSplitState()) {
                 BinlogOffset position = getBinlogPosition(element);
                 splitState.asBinlogSplitState().setStartingOffset(position);
@@ -141,6 +138,8 @@ public final class MySqlRecordEmitter<T>
             final TableChange tableSchema =
                     tableSchemas.getOrDefault(RecordUtils.getTableId(element), 
null);
 
+            updateSnapshotRecord(element, splitState);
+
             debeziumDeserializationSchema.deserialize(
                     element,
                     new Collector<T>() {
@@ -170,6 +169,12 @@ public final class MySqlRecordEmitter<T>
         }
     }
 
+    private void updateSnapshotRecord(SourceRecord element, MySqlSplitState 
splitState) {
+        if (splitState.isSnapshotSplitState() && includeIncremental) {
+            toSnapshotRecord(element);
+        }
+    }
+
     private void updateStartingOffsetForSplit(MySqlSplitState splitState, 
SourceRecord element) {
         if (splitState.isBinlogSplitState()) {
             // record the time metric to enter the incremental phase
diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/utils/RecordUtils.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/utils/RecordUtils.java
index 6944bd795..c05b58764 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/utils/RecordUtils.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/utils/RecordUtils.java
@@ -18,7 +18,9 @@
 package org.apache.inlong.sort.cdc.mysql.source.utils;
 
 import io.debezium.connector.AbstractSourceInfo;
+import io.debezium.connector.SnapshotRecord;
 import io.debezium.data.Envelope;
+import io.debezium.data.Envelope.FieldName;
 import io.debezium.document.DocumentReader;
 import io.debezium.relational.TableId;
 import io.debezium.relational.history.HistoryRecord;
@@ -247,6 +249,17 @@ public class RecordUtils {
         return watermarkKind.isPresent();
     }
 
+    public static boolean isSnapshotRecord(Struct sourceStruct) {
+        SnapshotRecord snapshotRecord = 
SnapshotRecord.fromSource(sourceStruct);
+        return (SnapshotRecord.TRUE == snapshotRecord);
+    }
+
+    public static void toSnapshotRecord(SourceRecord element) {
+        Struct messageStruct = (Struct) element.value();
+        Struct sourceStruct = messageStruct.getStruct(FieldName.SOURCE);
+        SnapshotRecord.TRUE.toSource(sourceStruct);
+    }
+
     public static boolean isLowWatermarkEvent(SourceRecord record) {
         Optional<SignalEventDispatcher.WatermarkKind> watermarkKind = 
getWatermarkKind(record);
         if (watermarkKind.isPresent() && watermarkKind.get() == 
SignalEventDispatcher.WatermarkKind.LOW) {
diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlReadableMetadata.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlReadableMetadata.java
index c5f339911..60c2c7e32 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlReadableMetadata.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlReadableMetadata.java
@@ -17,6 +17,8 @@
 
 package org.apache.inlong.sort.cdc.mysql.table;
 
+import static 
org.apache.inlong.sort.cdc.mysql.source.utils.RecordUtils.isSnapshotRecord;
+
 import io.debezium.connector.AbstractSourceInfo;
 import io.debezium.data.Envelope;
 import io.debezium.data.Envelope.FieldName;
@@ -173,7 +175,7 @@ public enum MySqlReadableMetadata {
                             .build();
                     DebeziumJson debeziumJson = 
DebeziumJson.builder().after(field).source(source)
                             
.tsMs(sourceStruct.getInt64(AbstractSourceInfo.TIMESTAMP_KEY)).op(getDebeziumOpType(data))
-                            .tableChange(tableSchema).build();
+                            
.tableChange(tableSchema).incremental(isSnapshotRecord(sourceStruct)).build();
 
                     try {
                         return 
StringData.fromString(OBJECT_MAPPER.writeValueAsString(debeziumJson));
@@ -453,7 +455,8 @@ public enum MySqlReadableMetadata {
                 .data(dataList).database(databaseName)
                 .sql("").es(opTs).isDdl(false).pkNames(getPkNames(tableSchema))
                 .mysqlType(getMysqlType(tableSchema)).table(tableName).ts(ts)
-                
.type(getCanalOpType(rowData)).sqlType(getSqlType(tableSchema)).build();
+                .type(getCanalOpType(rowData)).sqlType(getSqlType(tableSchema))
+                .incremental(isSnapshotRecord(sourceStruct)).build();
 
         try {
             return 
StringData.fromString(OBJECT_MAPPER.writeValueAsString(canalJson));
diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlTableInlongSourceFactory.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlTableInlongSourceFactory.java
index 6aaf06ad7..f9f78e2bb 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlTableInlongSourceFactory.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlTableInlongSourceFactory.java
@@ -48,6 +48,7 @@ import static 
org.apache.inlong.sort.cdc.mysql.source.config.MySqlSourceOptions.
 import static 
org.apache.inlong.sort.cdc.mysql.source.config.MySqlSourceOptions.DATABASE_NAME;
 import static 
org.apache.inlong.sort.cdc.mysql.source.config.MySqlSourceOptions.HEARTBEAT_INTERVAL;
 import static 
org.apache.inlong.sort.cdc.mysql.source.config.MySqlSourceOptions.HOSTNAME;
+import static 
org.apache.inlong.sort.cdc.mysql.source.config.MySqlSourceOptions.INCLUDE_INCREMENTAL;
 import static 
org.apache.inlong.sort.cdc.mysql.source.config.MySqlSourceOptions.MIGRATE_ALL;
 import static 
org.apache.inlong.sort.cdc.mysql.source.config.MySqlSourceOptions.PASSWORD;
 import static 
org.apache.inlong.sort.cdc.mysql.source.config.MySqlSourceOptions.PORT;
@@ -145,6 +146,7 @@ public class MySqlTableInlongSourceFactory implements 
DynamicTableSourceFactory
         int connectionPoolSize = config.get(CONNECTION_POOL_SIZE);
         final boolean appendSource = config.get(APPEND_MODE);
         final boolean migrateAll = config.get(MIGRATE_ALL);
+        final boolean includeIncremental = config.get(INCLUDE_INCREMENTAL);
         double distributionFactorUpper = 
config.get(SPLIT_KEY_EVEN_DISTRIBUTION_FACTOR_UPPER_BOUND);
         double distributionFactorLower = 
config.get(SPLIT_KEY_EVEN_DISTRIBUTION_FACTOR_LOWER_BOUND);
         boolean scanNewlyAddedTableEnabled = 
config.get(SCAN_NEWLY_ADDED_TABLE_ENABLED);
@@ -191,7 +193,9 @@ public class MySqlTableInlongSourceFactory implements 
DynamicTableSourceFactory
                 heartbeatInterval,
                 migrateAll,
                 inlongMetric,
-                inlongAudit, rowKindFiltered);
+                inlongAudit,
+                rowKindFiltered,
+                includeIncremental);
     }
 
     @Override
@@ -237,6 +241,7 @@ public class MySqlTableInlongSourceFactory implements 
DynamicTableSourceFactory
         options.add(INLONG_AUDIT);
         options.add(ROW_KINDS_FILTERED);
         options.add(AUDIT_KEYS);
+        options.add(INCLUDE_INCREMENTAL);
         return options;
     }
 
diff --git 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlTableSource.java
 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlTableSource.java
index d2effd978..fe283a04a 100644
--- 
a/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlTableSource.java
+++ 
b/inlong-sort/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/table/MySqlTableSource.java
@@ -84,6 +84,7 @@ public class MySqlTableSource implements ScanTableSource, 
SupportsReadingMetadat
     private final boolean migrateAll;
     private final String inlongMetric;
     private final String inlongAudit;
+    private final boolean includeIncremental;
     // 
--------------------------------------------------------------------------------------------
     // Mutable attributes
     // 
--------------------------------------------------------------------------------------------
@@ -127,7 +128,8 @@ public class MySqlTableSource implements ScanTableSource, 
SupportsReadingMetadat
             boolean migrateAll,
             String inlongMetric,
             String inlongAudit,
-            String rowKindsFiltered) {
+            String rowKindsFiltered,
+            boolean includeIncremental) {
         this(
                 physicalSchema,
                 port,
@@ -156,7 +158,8 @@ public class MySqlTableSource implements ScanTableSource, 
SupportsReadingMetadat
                 migrateAll,
                 inlongMetric,
                 inlongAudit,
-                rowKindsFiltered);
+                rowKindsFiltered,
+                includeIncremental);
     }
 
     /**
@@ -190,7 +193,8 @@ public class MySqlTableSource implements ScanTableSource, 
SupportsReadingMetadat
             boolean migrateAll,
             String inlongMetric,
             String inlongAudit,
-            String rowKindsFiltered) {
+            String rowKindsFiltered,
+            boolean includeIncremental) {
         this.physicalSchema = physicalSchema;
         this.port = port;
         this.hostname = checkNotNull(hostname);
@@ -222,6 +226,7 @@ public class MySqlTableSource implements ScanTableSource, 
SupportsReadingMetadat
         this.inlongMetric = inlongMetric;
         this.inlongAudit = inlongAudit;
         this.rowKindsFiltered = rowKindsFiltered;
+        this.includeIncremental = includeIncremental;
     }
 
     @Override
@@ -283,6 +288,7 @@ public class MySqlTableSource implements ScanTableSource, 
SupportsReadingMetadat
                             .heartbeatInterval(heartbeatInterval)
                             .inlongMetric(inlongMetric)
                             .inlongAudit(inlongAudit)
+                            .includeIncremental(includeIncremental)
                             .build();
             return SourceProvider.of(parallelSource);
         } else {
@@ -371,7 +377,9 @@ public class MySqlTableSource implements ScanTableSource, 
SupportsReadingMetadat
                         heartbeatInterval,
                         migrateAll,
                         inlongMetric,
-                        inlongAudit, rowKindsFiltered);
+                        inlongAudit,
+                        rowKindsFiltered,
+                        includeIncremental);
         source.metadataKeys = metadataKeys;
         source.producedDataType = producedDataType;
         return source;
diff --git 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateTest.java
 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateTest.java
index 0b5dfb6f4..e131971f3 100644
--- 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateTest.java
+++ 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateTest.java
@@ -54,6 +54,7 @@ public class AllMigrateTest {
         Map<String, String> option = new HashMap<>();
         option.put("append-mode", "true");
         option.put("migrate-all", "true");
+        option.put("include-incremental", "true");
         List<String> tables = new ArrayList(10);
         tables.add("test.*");
         List<FieldInfo> fields = Collections.singletonList(
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJson.java
 
b/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJson.java
index 5072c4e55..7ad548aa6 100644
--- 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJson.java
+++ 
b/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJson.java
@@ -20,127 +20,41 @@ package org.apache.inlong.sort.formats.json.canal;
 import java.util.List;
 import java.util.Map;
 import lombok.Builder;
+import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonInclude;
+import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonInclude.Include;
+import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
 
 @Builder
+@JsonInclude(Include.NON_NULL)
 public class CanalJson {
 
+    @JsonProperty("data")
     private List<Map<String, Object>> data;
+    @JsonProperty("es")
     private long es;
+    @JsonProperty("table")
     private String table;
+    @JsonProperty("type")
     private String type;
+    @JsonProperty("database")
     private String database;
+    @JsonProperty("ts")
     private long ts;
+    @JsonProperty("sql")
     private String sql;
+    @JsonProperty("mysqlType")
     private Map<String, String> mysqlType;
+    @JsonProperty("sqlType")
     private Map<String, Integer> sqlType;
-
+    @JsonProperty("isDdl")
     private boolean isDdl;
+    @JsonProperty("pkNames")
     private List<String> pkNames;
+    @JsonProperty("schema")
     private String schema;
+    @JsonProperty("oracleType")
     private Map<String, String> oracleType;
-
-    public List<Map<String, Object>> getData() {
-        return data;
-    }
-
-    public void setData(List<Map<String, Object>> data) {
-        this.data = data;
-    }
-
-    public long getEs() {
-        return es;
-    }
-
-    public void setEs(long es) {
-        this.es = es;
-    }
-
-    public String getTable() {
-        return table;
-    }
-
-    public void setTable(String table) {
-        this.table = table;
-    }
-
-    public String getType() {
-        return type;
-    }
-
-    public void setType(String type) {
-        this.type = type;
-    }
-
-    public String getDatabase() {
-        return database;
-    }
-
-    public void setDatabase(String database) {
-        this.database = database;
-    }
-
-    public long getTs() {
-        return ts;
-    }
-
-    public void setTs(long ts) {
-        this.ts = ts;
-    }
-
-    public String getSql() {
-        return sql;
-    }
-
-    public void setSql(String sql) {
-        this.sql = sql;
-    }
-
-    public Map<String, String> getMysqlType() {
-        return mysqlType;
-    }
-
-    public void setMysqlType(Map<String, String> mysqlType) {
-        this.mysqlType = mysqlType;
-    }
-
-    public void setDdl(boolean ddl) {
-        isDdl = ddl;
-    }
-
-    public List<String> getPkNames() {
-        return pkNames;
-    }
-
-    public void setPkNames(List<String> pkNames) {
-        this.pkNames = pkNames;
-    }
-
-    public Map<String, Integer> getSqlType() {
-        return sqlType;
-    }
-
-    public void setSqlType(Map<String, Integer> sqlType) {
-        this.sqlType = sqlType;
-    }
-
-    public boolean isDdl() {
-        return isDdl;
-    }
-
-    public Map<String, String> getOracleType() {
-        return oracleType;
-    }
-
-    public void setOracleType(Map<String, String> oracleType) {
-        this.oracleType = oracleType;
-    }
-
-    public String getSchema() {
-        return schema;
-    }
-
-    public void setSchema(String schema) {
-        this.schema = schema;
-    }
+    @JsonProperty("incremental")
+    private Boolean incremental;
 
 }
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJson.java
 
b/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJson.java
index f2e7cc9f3..feec23042 100644
--- 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJson.java
+++ 
b/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJson.java
@@ -33,6 +33,7 @@ public class DebeziumJson {
     private TableChanges.TableChange tableChange;
     private long tsMs;
     private String op;
+    private Boolean incremental;
 
     @Builder
     @Data

Reply via email to