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 =