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;
+        }
+    }
+}

Reply via email to