{"record":{"id":"98f7ee7d0375ba7c","repo":"apache/iceberg","slug":"malformed-record-topic-record-topic-partiti","errorCode":null,"errorMessage":"malformed record topic: ${record.topic()}, partition: ${record.kafkaPartition()}, offset: ${record.kafkaOffset()}","messagePattern":"malformed record topic: (.+?), partition: (.+?), offset: (.+?)","errorType":"exception","errorClass":"IllegalArgumentException","httpStatus":null,"severity":"error","filePath":"kafka-connect/kafka-connect-transforms/src/main/java/org/apache/iceberg/connect/transforms/MongoDebeziumTransform.java","lineNumber":113,"sourceCode":"          record.topic(),\n          record.kafkaPartition(),\n          null,\n          null,\n          null,\n          null,\n          record.timestamp(),\n          record.headers());\n    }\n\n    final SinkRecord keyIdRecord = keyIdExtractor.apply(record);\n    final SinkRecord afterRecord = afterExtractor.apply(record);\n    final SinkRecord beforeRecord = beforeExtractor.apply(record);\n    final SinkRecord updateDescriptionRecord = updateDescriptionExtractor.apply(record);\n\n    if (beforeRecord.value() == null\n        && afterRecord.value() == null\n        && updateDescriptionRecord.value() == null) {\n      throw new IllegalArgumentException(\n          String.format(\n              \"malformed record topic: %s, partition: %s, offset: %s\",\n              record.topic(), record.kafkaPartition(), record.kafkaOffset()));\n    }\n\n    BsonDocument afterBson = null;\n    BsonDocument beforeBson = null;\n    BsonDocument keyBson = BsonDocument.parse(\"{ \\\"id\\\" : \" + keyIdRecord.key().toString() + \"}\");\n\n    if (beforeRecord.value() != null) {\n      beforeBson = BsonDocument.parse(beforeRecord.value().toString());\n    }\n    if (afterRecord.value() == null && updateDescriptionRecord.value() != null) {\n      afterBson =\n          buildAfterBsonFromPartials(\n              updateDescriptionRecord,\n              (beforeBson == null) ? new BsonDocument() : beforeBson.clone(),\n              keyBson);","sourceCodeStart":95,"sourceCodeEnd":131,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/kafka-connect/kafka-connect-transforms/src/main/java/org/apache/iceberg/connect/transforms/MongoDebeziumTransform.java#L95-L131","documentation":"MongoDebeziumTransform.apply throws IllegalArgumentException when a Debezium MongoDB change event has null values for all of 'before', 'after', and 'updateDescription' — i.e. none of the expected payload sections are present, so the record cannot be interpreted as a change event.","triggerScenarios":"Applying the transform to a SinkRecord whose Debezium envelope lacks before, after, and updateDescription values simultaneously — e.g. tombstones, delete events without a 'before' image, or non-change-event records routed to the transform's topic.","commonSituations":"Misconfigured topics including tombstone records or heartbeat/DDL records reaching the transform; Debezium connector emitting delete events with extractors configured to skip 'before'; filter misconfiguration in the source connector.","solutions":["Configure the Debezium MongoDB connector to emit 'before' state (tombstones.on.delete, capture mode) so at least one section is populated.","Filter out tombstone/heartbeat records before the transform (e.g. DropTombstoneRecords SMT or topic routing).","Route only actual change-event topics to this transform."],"exampleFix":"// before — transform applied to all records including tombstones\n\"transforms\": \"mongo\"\n// after\n\"transforms\": \"dropTombstones,mongo\",\n\"transforms.dropTombstones.type\": \"org.apache.kafka.connect.transforms.DropTombstones$Value\"","handlingStrategy":"validation","validationCode":"// Guard before applying the transform\norg.apache.kafka.connect.data.Struct value = (Struct) record.value();\nboolean hasAny = value != null\n    && (value.schema().field(\"before\") != null || value.schema().field(\"after\") != null || value.schema().field(\"updateDescription\") != null);\n// route null-payload/tombstone records away from the transform","typeGuard":"boolean isChangeEvent(SinkRecord r) { return r.value() instanceof Struct s && (s.get(\"before\") != null || s.get(\"after\") != null || s.get(\"updateDescription\") != null); }","tryCatchPattern":"try { ... } catch (IllegalArgumentException e) { log.warn(\"Skipping non-change record: {}\", e.getMessage()); }","preventionTips":["Drop tombstones (DropTombstones SMT) before this transform","Set tombstones.on.delete=false on the Debezium source","Route only change-event topics to the transform"],"tags":["kafka-connect","debezium","mongodb"],"backgroundTag":"unexpected-response-shape","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}