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 =