apache/iceberg · error · DataException
An error occurred converting record, topic
Error message
An error occurred converting record, topic: ${record.topic()}, partition, ${record.kafkaPartition()}, offset: ${record.kafkaOffset()} What it means
IcebergWriter.write converts each SinkRecord's value into an Iceberg Record and writes it to the data file; any conversion or write failure for a single record is wrapped in a DataException annotated with the record's topic/partition/offset so the bad Kafka message can be located. The task will fail per Kafka Connect's error-reporting policy for that record.
Solutions
- Use the topic/partition/offset in the message to fetch the offending record and inspect its payload against the table schema.
- Align the record's schema with the Iceberg table — evolve the table via schema evolution or fix the producer.
- Configure Kafka Connect error handling (errors.tolerance=all, errors.deadletterqueue.topic.name) to quarantine bad records instead of failing the task.
- If values systematically mismatch, check SchemaUtils type inference settings and converter config (e.g. timestamp/decimal mappings).
Example fix
// before
// task fails on bad record
sink.connector.setProperties(Map.of(
"errors.tolerance", "none"));
// after
sink.connector.setProperties(Map.of(
"errors.tolerance", "all",
"errors.deadletterqueue.topic.name", "iceberg-sink-dlq",
"errors.log.enable", "true")); Defensive patterns
Strategy: try-catch
Validate before calling
// pre-flight: ensure record value/schema conform before write
if (record.value() == null || (record.valueSchema() == null && record.value().toString().isEmpty())) {
throw new DataException("Skipping unusable record: " + record.topic() + "/" + record.kafkaPartition() + "@" + record.kafkaOffset());
} Try / catch
try {
writer.write(record);
} catch (DataException e) {
LOG.error("Bad record at topic={} partition={} offset={}", record.topic(), record.kafkaPartition(), record.kafkaOffset(), e);
throw e; // or route to DLQ when errors.tolerance=all
} Prevention
- Configure errors.tolerance=all with a dead-letter queue on the sink connector
- Keep producer and Iceberg table schemas in sync via schema registry compatibility rules
- Validate numeric/timestamp ranges the Iceberg schema expects before producing
- Watch DLQ/alerts so schema drift is caught before mass task failure
When it happens
Trigger: convertToRow(record) throws (null/missing fields, type cast failures, schema vs payload mismatch, malformed nested values) or writer.write(row) throws (row does not match table schema, partition value out of bounds).
Common situations: Producer schema evolved (added/renamed fields) without updating the Iceberg table or value schema; null in a required column; numeric overflow when mapping Kafka types to Iceberg types; corrupted JSON/Avro payloads on the topic.
Understand the failure class
Background: Schema validation failed / invalid input schema: payload rejected because its shape doesn't match the expected schema — this error's family across 28 libraries.
Related errors
- An error occurred closing catalog instance, ignoring...
- Cannot convert date
- Cannot convert java.util.Date to variant without a…
- Cannot convert map to variant: keys must be non-null…
- Cannot convert Number to variant (unknown type)
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/d43fe0a141650fa2.
Report an issue: GitHub.
Appendix: source
Thrown at kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/IcebergWriter.java:67
this.writerResults = Lists.newArrayList();
initNewWriter();
}
private void initNewWriter() {
this.writer = RecordUtils.createTableWriter(table, tableReference, config);
this.recordConverter = new RecordConverter(table, config);
}
@Override
public void write(SinkRecord record) {
try {
// ignore tombstones...
if (record.value() != null) {
Record row = convertToRow(record);
writer.write(row);
}
} catch (Exception e) {
throw new DataException(
String.format(
Locale.ROOT,
"An error occurred converting record, topic: %s, partition, %d, offset: %d",
record.topic(),
record.kafkaPartition(),
record.kafkaOffset()),
e);
}
}
private Record convertToRow(SinkRecord record) {
if (!config.evolveSchemaEnabled()) {
return recordConverter.convert(record.value());
}
SchemaUpdate.Consumer updates = new SchemaUpdate.Consumer();
Record row = recordConverter.convert(record.value(), updates);
View on GitHub (pinned to 86d9c8fc54)