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

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


The following commit(s) were added to refs/heads/master by this push:
     new 1a3dcad6f2 [cdc] Fix NPE when metadata columns are used with 
debezium-bson format (#9055)
1a3dcad6f2 is described below

commit 1a3dcad6f2803851adfaa80dc859e59cc2b4c6b2
Author: Eunbin Son <[email protected]>
AuthorDate: Fri Aug 7 14:15:47 2026 +0900

    [cdc] Fix NPE when metadata columns are used with debezium-bson format 
(#9055)
---
 .../format/debezium/DebeziumBsonRecordParser.java  |  4 +++
 .../debezium/DebeziumBsonRecordParserTest.java     | 36 ++++++++++++++++++++++
 2 files changed, 40 insertions(+)

diff --git 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/format/debezium/DebeziumBsonRecordParser.java
 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/format/debezium/DebeziumBsonRecordParser.java
index 134ed8b383..23855d249f 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/format/debezium/DebeziumBsonRecordParser.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/format/debezium/DebeziumBsonRecordParser.java
@@ -109,6 +109,10 @@ public class DebeziumBsonRecordParser extends 
DebeziumJsonRecordParser {
 
     @Override
     protected void setRoot(CdcSourceRecord record) {
+        // Store current record for metadata access. Assign the field directly 
instead of calling
+        // super.setRoot, because DebeziumJsonRecordParser#setRoot also parses 
the Debezium value
+        // schema, which carries no field information for BSON documents.
+        this.currentRecord = record;
         root = (JsonNode) record.getValue();
         if (root.has(FIELD_SCHEMA)) {
             root = root.get(FIELD_PAYLOAD);
diff --git 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/format/debezium/DebeziumBsonRecordParserTest.java
 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/format/debezium/DebeziumBsonRecordParserTest.java
index 9c8dafc291..8c65753aa0 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/format/debezium/DebeziumBsonRecordParserTest.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/format/debezium/DebeziumBsonRecordParserTest.java
@@ -18,14 +18,17 @@
 
 package org.apache.paimon.flink.action.cdc.format.debezium;
 
+import org.apache.paimon.flink.action.cdc.CdcMetadataConverter;
 import org.apache.paimon.flink.action.cdc.CdcSourceRecord;
 import org.apache.paimon.flink.action.cdc.TypeMapping;
 import org.apache.paimon.flink.action.cdc.format.DataFormat;
+import org.apache.paimon.flink.action.cdc.kafka.KafkaMetadataConverter;
 import 
org.apache.paimon.flink.action.cdc.watermark.MessageQueueCdcTimestampExtractor;
 import org.apache.paimon.flink.sink.cdc.CdcRecord;
 import org.apache.paimon.flink.sink.cdc.CdcSchema;
 import org.apache.paimon.flink.sink.cdc.RichCdcMultiplexRecord;
 import org.apache.paimon.schema.Schema;
+import org.apache.paimon.types.DataField;
 import org.apache.paimon.types.RowKind;
 import org.apache.paimon.utils.JsonSerdeUtil;
 import org.apache.paimon.utils.StringUtils;
@@ -51,6 +54,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.stream.Collectors;
 
 /** Test for DebeziumBsonRecordParser. */
 public class DebeziumBsonRecordParserTest {
@@ -228,6 +232,38 @@ public class DebeziumBsonRecordParserTest {
         }
     }
 
+    @Test
+    public void extractRecordWithMetadataColumns() throws Exception {
+        DebeziumBsonRecordParser parser =
+                new DebeziumBsonRecordParser(TypeMapping.defaultMapping(), 
Collections.emptyList());
+        parser.withMetadataConverters(
+                new CdcMetadataConverter[] {
+                    new KafkaMetadataConverter.TopicConverter(),
+                    new KafkaMetadataConverter.OffsetConverter()
+                });
+
+        Assertions.assertFalse(insertList.isEmpty());
+        for (CdcSourceRecord cdcRecord : insertList) {
+            List<RichCdcMultiplexRecord> records = new ArrayList<>();
+            parser.flatMap(cdcRecord, new ListCollector<>(records));
+            Assertions.assertEquals(1, records.size());
+
+            Map<String, String> expected = new HashMap<>(beforeEvent);
+            expected.put("topic", "topic");
+            expected.put("offset", "0");
+
+            CdcRecord result = records.get(0).toRichCdcRecord().toCdcRecord();
+            Assertions.assertEquals(RowKind.INSERT, result.kind());
+            Assertions.assertEquals(expected, result.data());
+
+            List<String> fieldNames =
+                    records.get(0).buildSchema().fields().stream()
+                            .map(DataField::name)
+                            .collect(Collectors.toList());
+            
Assertions.assertTrue(fieldNames.containsAll(Arrays.asList("topic", "offset")));
+        }
+    }
+
     @Test
     public void bsonConvertJsonTest() throws Exception {
         DebeziumBsonRecordParser parser =

Reply via email to