From 0d1b187988882be990266966118ef9ecb9b26ef3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BC=A0=E4=B8=87=E4=B9=89?= Date: Fri, 21 Aug 2026 15:48:02 +0800 Subject: [PATCH] [Fix] MongoDB CDC: handle null fullDocument in transaction to prevent NullNode exception (#9332) When MongoDB processes a transaction with simultaneous update and delete operations, the Debezium CDC event may contain a null fullDocument for the update operation. Previously, getExtractRow() would call jsonNode.asText() on a NullNode, which returns the string "null", causing JsonSerdeUtil.asSpecificNodeType() to throw IllegalArgumentException because "null" parses to NullNode, not ObjectNode. This fix adds a null check at the beginning of getExtractRow() to gracefully handle null or NullNode jsonNode by returning an empty map. --- .../cdc/mongodb/strategy/Mongo4VersionStrategy.java | 13 +++++++++++-- .../cdc/mongodb/strategy/MongoVersionStrategy.java | 3 +++ 2 files changed, 14 insertions(+), 2 deletions(-) diff --git a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/Mongo4VersionStrategy.java b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/Mongo4VersionStrategy.java index 5f9538d2fcaf..bbdc7a8d018e 100644 --- a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/Mongo4VersionStrategy.java +++ b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/Mongo4VersionStrategy.java @@ -92,7 +92,10 @@ private List handleOperation( switch (op) { case OP_INSERT: - records.add(processRecord(fullDocument, RowKind.INSERT)); + RichCdcMultiplexRecord insertRecord = processRecord(fullDocument, RowKind.INSERT); + if (insertRecord != null) { + records.add(insertRecord); + } break; case OP_REPLACE: case OP_UPDATE: @@ -100,7 +103,10 @@ private List handleOperation( // information. Therefore, data is first deleted using the primary key '_id', and // then inserted. records.add(processRecord(documentKey, RowKind.DELETE)); - records.add(processRecord(fullDocument, RowKind.INSERT)); + RichCdcMultiplexRecord updateRecord = processRecord(fullDocument, RowKind.INSERT); + if (updateRecord != null) { + records.add(updateRecord); + } break; case OP_DELETE: records.add(processRecord(documentKey, RowKind.DELETE)); @@ -125,6 +131,9 @@ private RichCdcMultiplexRecord processRecord(JsonNode fullDocument, RowKind rowK CdcSchema.Builder schemaBuilder = CdcSchema.newBuilder(); Map record = getExtractRow(fullDocument, schemaBuilder, computedColumns, mongodbConfig); + if (record == null) { + return null; + } schemaBuilder.primaryKey(extractPrimaryKeys()); return new RichCdcMultiplexRecord( databaseName, collection, schemaBuilder.build(), new CdcRecord(rowKind, record)); diff --git a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/MongoVersionStrategy.java b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/MongoVersionStrategy.java index a41a1eec6874..fcabcaffbba7 100644 --- a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/MongoVersionStrategy.java +++ b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/MongoVersionStrategy.java @@ -82,6 +82,9 @@ default Map getExtractRow( List computedColumns, Configuration mongodbConfig) throws JsonProcessingException { + if (jsonNode == null || jsonNode.isNull()) { + return null; + } SchemaAcquisitionMode mode = SchemaAcquisitionMode.valueOf(mongodbConfig.get(START_MODE).toUpperCase()); ObjectNode objectNode =