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