{"record":{"id":"d43fe0a141650fa2","repo":"apache/iceberg","slug":"an-error-occurred-converting-record-topic-reco","errorCode":null,"errorMessage":"An error occurred converting record, topic: ${record.topic()}, partition, ${record.kafkaPartition()}, offset: ${record.kafkaOffset()}","messagePattern":"An error occurred converting record, topic: (.+?), partition, (.+?), offset: (.+?)","errorType":"exception","errorClass":"DataException","httpStatus":null,"severity":"error","filePath":"kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/IcebergWriter.java","lineNumber":67,"sourceCode":"    this.writerResults = Lists.newArrayList();\n    initNewWriter();\n  }\n\n  private void initNewWriter() {\n    this.writer = RecordUtils.createTableWriter(table, tableReference, config);\n    this.recordConverter = new RecordConverter(table, config);\n  }\n\n  @Override\n  public void write(SinkRecord record) {\n    try {\n      // ignore tombstones...\n      if (record.value() != null) {\n        Record row = convertToRow(record);\n        writer.write(row);\n      }\n    } catch (Exception e) {\n      throw new DataException(\n          String.format(\n              Locale.ROOT,\n              \"An error occurred converting record, topic: %s, partition, %d, offset: %d\",\n              record.topic(),\n              record.kafkaPartition(),\n              record.kafkaOffset()),\n          e);\n    }\n  }\n\n  private Record convertToRow(SinkRecord record) {\n    if (!config.evolveSchemaEnabled()) {\n      return recordConverter.convert(record.value());\n    }\n\n    SchemaUpdate.Consumer updates = new SchemaUpdate.Consumer();\n    Record row = recordConverter.convert(record.value(), updates);\n","sourceCodeStart":49,"sourceCodeEnd":85,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/IcebergWriter.java#L49-L85","documentation":"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.","triggerScenarios":"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).","commonSituations":"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.","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)."],"exampleFix":"// before\n// task fails on bad record\nsink.connector.setProperties(Map.of(\n    \"errors.tolerance\", \"none\"));\n\n// after\nsink.connector.setProperties(Map.of(\n    \"errors.tolerance\", \"all\",\n    \"errors.deadletterqueue.topic.name\", \"iceberg-sink-dlq\",\n    \"errors.log.enable\", \"true\"));","handlingStrategy":"try-catch","validationCode":"// pre-flight: ensure record value/schema conform before write\nif (record.value() == null || (record.valueSchema() == null && record.value().toString().isEmpty())) {\n  throw new DataException(\"Skipping unusable record: \" + record.topic() + \"/\" + record.kafkaPartition() + \"@\" + record.kafkaOffset());\n}","typeGuard":null,"tryCatchPattern":"try {\n  writer.write(record);\n} catch (DataException e) {\n  LOG.error(\"Bad record at topic={} partition={} offset={}\", record.topic(), record.kafkaPartition(), record.kafkaOffset(), e);\n  throw e; // or route to DLQ when errors.tolerance=all\n}","preventionTips":["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"],"tags":["kafka-connect","record-conversion","schema-mismatch","data-exception"],"backgroundTag":"schema-validation-failed","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"}