This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 6472c5ec59 [Fix][Connector-V2] Fix TiDB CDC delete deserialization
(#11223)
6472c5ec59 is described below
commit 6472c5ec59b62552d9b6931889d3523d470df905
Author: zhaoysg <[email protected]>
AuthorDate: Thu Aug 20 22:50:41 2026 +0800
[Fix][Connector-V2] Fix TiDB CDC delete deserialization (#11223)
Co-authored-by: davidzollo <[email protected]>
Co-authored-by: DanielLeens <[email protected]>
Co-authored-by: davidzollo <[email protected]>
---
.../AbstractSeaTunnelRowDeserializer.java | 9 +
.../SeaTunnelRowSnapshotRecordDeserializer.java | 2 +-
.../SeaTunnelRowStreamingRecordDeserializer.java | 42 +++-
...eaTunnelRowStreamingRecordDeserializerTest.java | 227 +++++++++++++++++++++
4 files changed, 275 insertions(+), 5 deletions(-)
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/AbstractSeaTunnelRowDeserializer.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/AbstractSeaTunnelRowDeserializer.java
index 77654239d1..b41457d57d 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/AbstractSeaTunnelRowDeserializer.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/AbstractSeaTunnelRowDeserializer.java
@@ -35,5 +35,14 @@ public abstract class
AbstractSeaTunnelRowDeserializer<Input> {
this.catalogTable = catalogTable;
}
+ /**
+ * Attach the source table identity before emitting rows so downstream
multi-table sinks can
+ * route TiDB CDC records by table.
+ */
+ protected void collect(SeaTunnelRow row, Collector<SeaTunnelRow> output) {
+ row.setTableId(catalogTable.getTablePath().getFullName());
+ output.collect(row);
+ }
+
abstract void deserialize(Input record, Collector<SeaTunnelRow> output)
throws Exception;
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowSnapshotRecordDeserializer.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowSnapshotRecordDeserializer.java
index e7c7fc1812..d6d398e8a3 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowSnapshotRecordDeserializer.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowSnapshotRecordDeserializer.java
@@ -53,6 +53,6 @@ public class SeaTunnelRowSnapshotRecordDeserializer
RowKey.decode(record.getKey().toByteArray()).getHandle(),
tableInfo);
SeaTunnelRow row = converter.convert(values, tableInfo, rowType);
- output.collect(row);
+ collect(row, output);
}
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowStreamingRecordDeserializer.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowStreamingRecordDeserializer.java
index 5e286e591e..0bc5ea7e35 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowStreamingRecordDeserializer.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowStreamingRecordDeserializer.java
@@ -25,8 +25,10 @@ import
org.apache.seatunnel.connectors.seatunnel.cdc.tidb.source.converter.DataC
import
org.apache.seatunnel.connectors.seatunnel.cdc.tidb.source.converter.DefaultDataConverter;
import org.tikv.common.key.RowKey;
+import org.tikv.common.meta.TiColumnInfo;
import org.tikv.common.meta.TiTableInfo;
import org.tikv.kvproto.Cdcpb;
+import org.tikv.shade.com.google.protobuf.ByteString;
import static org.tikv.common.codec.TableCodec.decodeObjects;
@@ -49,10 +51,10 @@ public class SeaTunnelRowStreamingRecordDeserializer
Object[] values;
switch (row.getOpType()) {
case DELETE:
- values = decodeObjects(row.getOldValue().toByteArray(),
handle, tableInfo);
+ values = decodeDeleteValues(row, handle);
SeaTunnelRow record = converter.convert(values, tableInfo,
rowType);
record.setRowKind(RowKind.DELETE);
- output.collect(record);
+ collect(record, output);
break;
case PUT:
try {
@@ -64,11 +66,11 @@ public class SeaTunnelRowStreamingRecordDeserializer
if (row.getOldValue() == null ||
row.getOldValue().isEmpty()) {
SeaTunnelRow insert = converter.convert(values,
tableInfo, rowType);
insert.setRowKind(RowKind.INSERT);
- output.collect(insert);
+ collect(insert, output);
} else {
SeaTunnelRow update = converter.convert(values,
tableInfo, rowType);
update.setRowKind(RowKind.UPDATE_AFTER);
- output.collect(update);
+ collect(update, output);
}
break;
} catch (final RuntimeException e) {
@@ -82,4 +84,36 @@ public class SeaTunnelRowStreamingRecordDeserializer
throw new IllegalArgumentException("Unknown Row Op Type: " +
row.getOpType());
}
}
+
+ private Object[] decodeDeleteValues(Cdcpb.Event.Row row, long handle) {
+ ByteString oldValue = row.getOldValue();
+ if (oldValue != null && !oldValue.isEmpty()) {
+ return decodeObjects(oldValue.toByteArray(), handle, tableInfo);
+ }
+ ByteString value = row.getValue();
+ if (value != null && !value.isEmpty()) {
+ // Prefer an available row image before falling back to a PK-only
delete row.
+ return decodeObjects(value.toByteArray(), handle, tableInfo);
+ }
+ return decodeDeleteValuesFromHandle(row, handle);
+ }
+
+ private Object[] decodeDeleteValuesFromHandle(Cdcpb.Event.Row row, long
handle) {
+ TiColumnInfo handleColumn = tableInfo.getPKIsHandleColumn();
+ if (!tableInfo.isPkHandle() || handleColumn == null) {
+ throw new IllegalStateException(
+ String.format(
+ "Cannot deserialize TiDB DELETE CDC row because
both value and"
+ + " oldValue are empty, and table %s(%s)
is not a pk-handle"
+ + " table. key=%s, startTs=%s,
commitTs=%s",
+ tableInfo.getName(),
+ tableInfo.getId(),
+ row.getKey(),
+ row.getStartTs(),
+ row.getCommitTs()));
+ }
+ Object[] values = new Object[tableInfo.getColumns().size()];
+ values[handleColumn.getOffset()] = handle;
+ return values;
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowStreamingRecordDeserializerTest.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowStreamingRecordDeserializerTest.java
new file mode 100644
index 0000000000..4fbc4ced33
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/deserializer/SeaTunnelRowStreamingRecordDeserializerTest.java
@@ -0,0 +1,227 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.cdc.tidb.source.deserializer;
+
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.PhysicalColumn;
+import org.apache.seatunnel.api.table.catalog.TableIdentifier;
+import org.apache.seatunnel.api.table.catalog.TableSchema;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.RowKind;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+
+import org.junit.jupiter.api.Test;
+import org.tikv.common.codec.TableCodec;
+import org.tikv.common.key.RowKey;
+import org.tikv.common.meta.CIStr;
+import org.tikv.common.meta.TiColumnInfo;
+import org.tikv.common.meta.TiTableInfo;
+import org.tikv.common.types.IntegerType;
+import org.tikv.common.types.StringType;
+import org.tikv.kvproto.Cdcpb;
+import org.tikv.kvproto.Kvrpcpb;
+import org.tikv.shade.com.google.protobuf.ByteString;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class SeaTunnelRowStreamingRecordDeserializerTest {
+
+ private static final long TABLE_ID = 42L;
+ private static final long HANDLE = 7L;
+
+ @Test
+ void deserializeDeleteShouldUseValueWhenOldValueIsEmpty() throws Exception
{
+ TiTableInfo tableInfo = tableInfo(true);
+ SeaTunnelRowStreamingRecordDeserializer deserializer =
+ new SeaTunnelRowStreamingRecordDeserializer(tableInfo,
catalogTable());
+ byte[] encodedRow =
+ TableCodec.encodeRow(
+ tableInfo.getColumns(), new Object[] {HANDLE,
"Alice"}, true, false);
+ TestCollector collector = new TestCollector();
+
+ deserializer.deserialize(deleteRow(ByteString.copyFrom(encodedRow)),
collector);
+
+ SeaTunnelRow row = collector.rows.get(0);
+ assertEquals(RowKind.DELETE, row.getRowKind());
+ assertEquals("test_db.test_table", row.getTableId());
+ assertEquals(HANDLE, row.getField(0));
+ assertEquals("Alice", row.getField(1));
+ }
+
+ @Test
+ void deserializeDeleteShouldUseHandleWhenValuesAreEmptyForPkHandleTable()
throws Exception {
+ SeaTunnelRowStreamingRecordDeserializer deserializer =
+ new SeaTunnelRowStreamingRecordDeserializer(tableInfo(true),
catalogTable());
+ TestCollector collector = new TestCollector();
+
+ deserializer.deserialize(deleteRow(ByteString.EMPTY), collector);
+
+ SeaTunnelRow row = collector.rows.get(0);
+ assertEquals(RowKind.DELETE, row.getRowKind());
+ assertEquals("test_db.test_table", row.getTableId());
+ assertEquals(HANDLE, row.getField(0));
+ assertNull(row.getField(1));
+ }
+
+ @Test
+ void deserializePutShouldSetTableId() throws Exception {
+ TiTableInfo tableInfo = tableInfo(true);
+ SeaTunnelRowStreamingRecordDeserializer deserializer =
+ new SeaTunnelRowStreamingRecordDeserializer(tableInfo,
catalogTable());
+ byte[] encodedRow =
+ TableCodec.encodeRow(
+ tableInfo.getColumns(), new Object[] {HANDLE,
"Alice"}, true, false);
+ TestCollector collector = new TestCollector();
+
+ deserializer.deserialize(putRow(ByteString.copyFrom(encodedRow)),
collector);
+
+ SeaTunnelRow row = collector.rows.get(0);
+ assertEquals(RowKind.INSERT, row.getRowKind());
+ assertEquals("test_db.test_table", row.getTableId());
+ assertEquals(HANDLE, row.getField(0));
+ assertEquals("Alice", row.getField(1));
+ }
+
+ @Test
+ void deserializeSnapshotShouldSetTableId() throws Exception {
+ TiTableInfo tableInfo = tableInfo(true);
+ SeaTunnelRowSnapshotRecordDeserializer deserializer =
+ new SeaTunnelRowSnapshotRecordDeserializer(tableInfo,
catalogTable());
+ byte[] encodedRow =
+ TableCodec.encodeRow(
+ tableInfo.getColumns(), new Object[] {HANDLE,
"Alice"}, true, false);
+ TestCollector collector = new TestCollector();
+
+ deserializer.deserialize(
+ Kvrpcpb.KvPair.newBuilder()
+ .setKey(RowKey.toRowKey(TABLE_ID,
HANDLE).toByteString())
+ .setValue(ByteString.copyFrom(encodedRow))
+ .build(),
+ collector);
+
+ SeaTunnelRow row = collector.rows.get(0);
+ assertEquals("test_db.test_table", row.getTableId());
+ assertEquals(HANDLE, row.getField(0));
+ assertEquals("Alice", row.getField(1));
+ }
+
+ @Test
+ void
deserializeDeleteShouldFailClearlyWhenValuesAreEmptyForNonPkHandleTable() {
+ SeaTunnelRowStreamingRecordDeserializer deserializer =
+ new SeaTunnelRowStreamingRecordDeserializer(tableInfo(false),
catalogTable());
+ TestCollector collector = new TestCollector();
+
+ IllegalStateException exception =
+ assertThrows(
+ IllegalStateException.class,
+ () ->
deserializer.deserialize(deleteRow(ByteString.EMPTY), collector));
+
+ assertTrue(exception.getMessage().contains("both value and oldValue
are empty"));
+ }
+
+ private static Cdcpb.Event.Row deleteRow(ByteString value) {
+ return Cdcpb.Event.Row.newBuilder()
+ .setType(Cdcpb.Event.LogType.PREWRITE)
+ .setOpType(Cdcpb.Event.Row.OpType.DELETE)
+ .setKey(RowKey.toRowKey(TABLE_ID, HANDLE).toByteString())
+ .setValue(value)
+ .build();
+ }
+
+ private static Cdcpb.Event.Row putRow(ByteString value) {
+ return Cdcpb.Event.Row.newBuilder()
+ .setType(Cdcpb.Event.LogType.PREWRITE)
+ .setOpType(Cdcpb.Event.Row.OpType.PUT)
+ .setKey(RowKey.toRowKey(TABLE_ID, HANDLE).toByteString())
+ .setValue(value)
+ .build();
+ }
+
+ private static TiTableInfo tableInfo(boolean pkIsHandle) {
+ List<TiColumnInfo> columns =
+ Arrays.asList(
+ new TiColumnInfo(1L, "id", 0, IntegerType.BIGINT,
true),
+ new TiColumnInfo(2L, "name", 1, StringType.VARCHAR,
false));
+ return new TiTableInfo(
+ TABLE_ID,
+ CIStr.newCIStr("test_table"),
+ "utf8mb4",
+ "utf8mb4_bin",
+ pkIsHandle,
+ columns,
+ Collections.emptyList(),
+ "",
+ 0L,
+ 2L,
+ 0L,
+ 0L,
+ null,
+ null,
+ null,
+ 0L,
+ 0L,
+ 0L,
+ null);
+ }
+
+ private static CatalogTable catalogTable() {
+ TableSchema tableSchema =
+ TableSchema.builder()
+ .column(
+ PhysicalColumn.of(
+ "id", BasicType.LONG_TYPE, (Long)
null, false, null, null))
+ .column(
+ PhysicalColumn.of(
+ "name",
+ BasicType.STRING_TYPE,
+ (Long) null,
+ true,
+ null,
+ null))
+ .build();
+ return CatalogTable.of(
+ TableIdentifier.of("test_catalog", "test_db", "test_table"),
+ tableSchema,
+ Collections.emptyMap(),
+ Collections.emptyList(),
+ null);
+ }
+
+ private static class TestCollector implements Collector<SeaTunnelRow> {
+ private final List<SeaTunnelRow> rows = new ArrayList<>();
+
+ @Override
+ public void collect(SeaTunnelRow record) {
+ rows.add(record);
+ }
+
+ @Override
+ public Object getCheckpointLock() {
+ return this;
+ }
+ }
+}