From 2c9ca47859c39d1f5870be119f133d024744e4f6 Mon Sep 17 00:00:00 2001 From: Eunbin Son Date: Thu, 6 Aug 2026 10:31:38 +0900 Subject: [PATCH] [cdc] Fix NPE when metadata columns are used with debezium-bson format DebeziumBsonRecordParser#setRoot did not store the current record, so AbstractRecordParser#evalMetadataColumns passed null to the metadata converters and threw NPE on the first record whenever --metadata_column was used. The sibling parsers already store it (#7315). Generated-by: Claude Code --- .../debezium/DebeziumBsonRecordParser.java | 4 +++ .../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 134ed8b3831c..23855d249fa5 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 List extractRecords() { @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 9c8dafc291a3..8c65753aa0bc 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.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 void extractDeleteRecord() throws Exception { } } + @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 records = new ArrayList<>(); + parser.flatMap(cdcRecord, new ListCollector<>(records)); + Assertions.assertEquals(1, records.size()); + + Map 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 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 =