apache/iceberg · error · IllegalArgumentException
malformed record topic
Error message
malformed record topic: ${record.topic()}, partition: ${record.kafkaPartition()}, offset: ${record.kafkaOffset()} What it means
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.
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.
Example fix
// before — transform applied to all records including tombstones "transforms": "mongo" // after "transforms": "dropTombstones,mongo", "transforms.dropTombstones.type": "org.apache.kafka.connect.transforms.DropTombstones$Value"
Defensive patterns
Strategy: validation
Validate before calling
// Guard before applying the transform
org.apache.kafka.connect.data.Struct value = (Struct) record.value();
boolean hasAny = value != null
&& (value.schema().field("before") != null || value.schema().field("after") != null || value.schema().field("updateDescription") != null);
// route null-payload/tombstone records away from the transform Type guard
boolean isChangeEvent(SinkRecord r) { return r.value() instanceof Struct s && (s.get("before") != null || s.get("after") != null || s.get("updateDescription") != null); } Try / catch
try { ... } catch (IllegalArgumentException e) { log.warn("Skipping non-change record: {}", e.getMessage()); } Prevention
- Drop tombstones (DropTombstones SMT) before this transform
- Set tombstones.on.delete=false on the Debezium source
- Route only change-event topics to the transform
When it happens
Trigger: 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.
Common situations: 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.
Understand the failure class
- Parsing and encoding errors: unexpected token, malformed input — why parsers reject input and how to find the real culprit.
Related errors
- Failed to find field
- Field of schema is not the same type for all documents in…
- Field of schema is not a homogenous array. Check option…
- The value type ' is not yet supported inside for a…
- An error occurred closing catalog instance, ignoring...
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/98f7ee7d0375ba7c.
Report an issue: GitHub.
Appendix: source
Thrown at kafka-connect/kafka-connect-transforms/src/main/java/org/apache/iceberg/connect/transforms/MongoDebeziumTransform.java:113
record.topic(),
record.kafkaPartition(),
null,
null,
null,
null,
record.timestamp(),
record.headers());
}
final SinkRecord keyIdRecord = keyIdExtractor.apply(record);
final SinkRecord afterRecord = afterExtractor.apply(record);
final SinkRecord beforeRecord = beforeExtractor.apply(record);
final SinkRecord updateDescriptionRecord = updateDescriptionExtractor.apply(record);
if (beforeRecord.value() == null
&& afterRecord.value() == null
&& updateDescriptionRecord.value() == null) {
throw new IllegalArgumentException(
String.format(
"malformed record topic: %s, partition: %s, offset: %s",
record.topic(), record.kafkaPartition(), record.kafkaOffset()));
}
BsonDocument afterBson = null;
BsonDocument beforeBson = null;
BsonDocument keyBson = BsonDocument.parse("{ \"id\" : " + keyIdRecord.key().toString() + "}");
if (beforeRecord.value() != null) {
beforeBson = BsonDocument.parse(beforeRecord.value().toString());
}
if (afterRecord.value() == null && updateDescriptionRecord.value() != null) {
afterBson =
buildAfterBsonFromPartials(
updateDescriptionRecord,
(beforeBson == null) ? new BsonDocument() : beforeBson.clone(),
keyBson);View on GitHub (pinned to 86d9c8fc54)