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

  1. Use the topic/partition/offset in the message to fetch the offending record and inspect its payload against the table schema.
  2. Align the record's schema with the Iceberg table — evolve the table via schema evolution or fix the producer.
  3. Configure Kafka Connect error handling (errors.tolerance=all, errors.deadletterqueue.topic.name) to quarantine bad records instead of failing the task.
  4. 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

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


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)