Skip to content

Kafka to Iceberg pipeline breaks after a schema change

Iceberg
Chad Harris·October 3, 2026·12 min read

A producer team ships a new version of an event. Within minutes the Kafka Connect Iceberg sink shows a task in the FAILED state, or the Flink job writing the table starts restarting. In the quieter version nothing alerts. The connector stays RUNNING and rows keep landing, while a column the producer renamed stops filling in or a field the producer added never reaches the table. The Apache Iceberg project describes engines such as Spark, Trino, Flink, Presto, Hive and Impala working with the same tables at the same time, so a gap in the table shows up in every engine that reads it.

A schema change failure in a Kafka to Iceberg pipeline is a mismatch between the schema the producer now writes and the schema of the Iceberg table. The writer either fails on the mismatch or writes the rows with data missing. The Apache Iceberg Kafka Connect sink is the self-managed route from a Kafka topic into an Iceberg table, and its feature list includes automatic table creation and schema evolution, both switched off by default.

This page covers the schema step. If the writer has the right schema but its commits fail with CommitFailedException, see Iceberg commits failing or conflicting. If data lands but the table holds thousands of tiny files, see too many small files and how to compact them. For how tables, snapshots and schemas fit together, see the complete Iceberg guide.

Check whether a task failed

Start by finding out whether anything failed at all. The Kafka Connect REST API has a status endpoint that reports whether the connector is running, failed or paused, the error information if it has failed, and the state of every task.

GET /connectors/{name}/status

For a task in the FAILED state, the error information includes the stack trace. Read the first exception in it. A deserialization error from the converter points at the record framing or the registry, covered below. A CommitFailedException or a commit timeout is a commit problem, not a schema one, and the commit failures page applies. A broader walk through Connect failures is in the Kafka Connect troubleshooting guide.

Check whether fields are silently dropped

If every task is RUNNING, the pipeline may be dropping data rather than failing. Compare the newest record on the topic with the newest rows in the table. A field present in the record and missing from the table schema, or a column that is null in every row written after the producer’s release, is the silent version of this problem. The Iceberg metadata tables list each snapshot with its committed_at time, so the first commit after the release is easy to find.

SELECT committed_at, snapshot_id, summary
FROM prod.db.events.snapshots
ORDER BY committed_at DESC;

Check how the records are framed

Before blaming the schema change itself, confirm the sink can still read the records. A registry-aware deserializer has to find the schema ID before it can decode anything. The Apicurio Registry SerDes documentation describes the ID location being found by checking for a magic byte at the start of the payload, and reading the ID from the message headers if the magic byte is not there. A producer that switched serializer, or started writing plain JSON or Protobuf with no registry, breaks the converter on the first new record, whatever the schema says. The Kafka deserialization error guide covers that failure from the consumer’s side.

The AWS Glue Schema Registry documentation says a schema version can be referenced by its UUID or version number, and Glue ships its own SerDes libraries. Google’s managed registry, by contrast, implements the Confluent Schema Registry REST API, so the same client libraries work against it. A converter has to match the framing the producer used, so a producer that moves between registries can look like a schema change from the sink’s side.

The converter the sink runs is set on the connector, for example:

value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=https://registry.example.com

What Iceberg allows

The Iceberg evolution documentation lists five schema changes and says they are metadata changes, so no data files are rewritten.

  • Add a column.
  • Drop a column.
  • Rename a column.
  • Update a column, which widens its type.
  • Reorder columns.

Iceberg uses unique IDs to track each column, and an added column gets a new ID, so existing data is never used by mistake.

The Iceberg table spec sets the detail. Renaming a field changes its name but not its field ID. Deleting a field removes it from the current schema, and the deletion cannot be rolled back unless the field was nullable or the current snapshot has not changed. The valid type promotions are:

int to long
float to double
decimal(P,S) to decimal(P',S), where P' > P
date to timestamp or timestamp_ns, in format v3 and later

Any other type change, such as int to string, is not a promotion. The spec also blocks a promotion on a field that a partition transform uses if the transform would produce a different value afterwards.

What the Kafka Connect sink does

The sink’s behaviour on a schema change is set by four properties, all of them false by default according to the sink configuration reference.

iceberg.tables.auto-create-enabled=false
iceberg.tables.evolve-schema-enabled=false
iceberg.tables.schema-force-optional=false
iceberg.tables.schema-case-insensitive=false

The documentation describes evolve-schema-enabled as adding any missing record fields to the table schema. What happens when it is off is visible in the sink’s record converter. The source at release 1.12.0 looks each record field up in the table schema by name. When there is no matching column, it adds one if schema evolution is on and otherwise skips the value. A new field from the producer therefore lands nowhere, and the task stays RUNNING.

With evolution on, the same converter adds missing columns, makes a required column optional when the record field is optional, and widens a column in two cases, which the schema utilities limit to int to long and float to double. It has no update for dropping or renaming a column. A field the producer removes leaves its column in place, with no value in new rows. A field the producer renames is a name the sink does not know, so it becomes a new column, and the old column stops receiving data.

Figure F1 sets each producer change beside what Iceberg allows and what the sink does, using the evolution documentation and the sink source linked above.

F1 Each producer change, in Iceberg and in the Kafka Connect sink
What Iceberg allows Kafka Connect sink with evolve-schema-enabled=true
New field New column, new field ID. Column added to the table.
Renamed field Name changes, field ID stays. New name added as a new column, unless a name mapping lists both names.
Removed field Column dropped, metadata only. Column stays. New rows carry no value.
Widened type Listed promotions only. int to long and float to double.
Other type change Not a valid promotion. No type update made.
Required field made optional Allowed. Column made optional.
Sources: the Iceberg evolution docs and the Kafka Connect sink source at release 1.12.0, both linked in the text, and the Iceberg table spec

The Iceberg Flink writes documentation describes schema evolution for the Dynamic Iceberg Sink. It tries to match each record’s schema to an existing table schema, adapts the record if it can (for example by adding a null for an extra optional table column), and otherwise evolves the table. The supported updates are adding columns, widening column types such as Integer to Long and Float to Double, and making required columns optional.

Dropping columns is disabled by default and is turned on with the builder option below. The documentation warns that a dropped field that reappears comes back as an entirely new column, so queries for it never return the old column’s data.

dropUnusedColumns(true)

Renaming is unsupported because the Dynamic Sink compares schemas by name. For a job written against a fixed schema with the standard sink, the table schema and the job’s row type are set in code, so a producer change needs a code change and a redeploy of the job.

What happens when the partition spec changes

A team can also change how a table is partitioned while the pipeline runs. The Iceberg evolution documentation says partitioning can be updated on an existing table because queries do not reference partition values directly. Data written under an earlier spec stays in its old layout, new data is written with the new spec, and the change is a metadata operation that rewrites no files. Readers are not affected, because Iceberg uses hidden partitioning and prunes files without a query naming the layout. On Spark the DDL documentation gives the statement.

ALTER TABLE prod.db.events ADD PARTITION FIELD year(ts);

The Kafka Connect sink does not change a table’s spec. Its documented iceberg.tables.default-partition-by setting is a list of partition fields to use when creating tables, and the sink’s table creation code is the only place the sink applies it. A spec change is therefore made on the table, not in the connector configuration. The sink documentation does not say when a running connector picks up a changed spec, so after the change, check the spec_id column of the table’s files metadata table to confirm new files carry the new spec.

How registry compatibility lets a breaking change through

A schema registry checks each new version of a subject against its compatibility mode before accepting it. The Confluent Schema Registry configuration source sets backward as the default for schema.compatibility.level. The AWS Glue Schema Registry documentation lists eight modes and describes BACKWARD as the check to use when deleting fields or adding optional fields.

NONE, DISABLED, BACKWARD, BACKWARD_ALL, FORWARD, FORWARD_ALL, FULL, FULL_ALL

Google’s registry overview defines backward compatibility as a consumer with the new schema version being able to read data produced with a previous one. That is a guarantee about reading Kafka records. It says nothing about the Iceberg table. Three changes pass a BACKWARD check and still change what lands in the table:

  • Adding an optional field. It passes, and with evolve-schema-enabled=false the sink skips its values.
  • Deleting a field. It passes, and the column stays in the table with no values in new rows.
  • Renaming a field by removing the old name and adding the new one as optional. It passes, and both sinks treat it as a new column, so the data splits across two columns at the release time.

Under NONE, nothing is checked, and a type change such as long to string reaches the sink. The Iceberg table cannot take it as a promotion, so the fix is a new column or a new table, chosen before the producer ships.

The subject’s mode is read from the registry API.

GET /config/{subject}

A change to the mode is a governance decision, so it belongs with whoever owns the subject. The schema registry tools guide compares the interfaces for checking and changing it, and separate comparisons cover Confluent Schema Registry and AWS Glue Schema Registry.

Fix it, one change at a time

Restart only the failed tasks

If the task failed on a record it could not convert, fix the cause first, then restart. The Kafka Connect REST API restarts the connector and only the tasks that failed with one call.

POST /connectors/{name}/restart?includeTasks=true&onlyFailed=true

Restarting without a fix replays the same record and fails again. For records that can never convert, the Connect error handling settings route them to a dead letter queue topic rather than failing the task. The Connect user guide lists them.

errors.tolerance=all
errors.deadletterqueue.topic.name=iceberg-sink-dlq

A dead letter queue keeps the pipeline moving, and every record in it is a row missing from the table until it is replayed.

Turn on schema evolution for additive changes

If producers mostly add optional fields, letting the sink add columns removes the silent drop.

iceberg.tables.evolve-schema-enabled=true

Turn it on before the producer’s next release. It does not bring back values skipped before it was enabled, which need the backfill below.

Make new columns optional

If the sink creates or evolves tables from records whose fields are marked required, later records without those fields cannot be written. The sink can create every column as optional instead.

iceberg.tables.schema-force-optional=true

Rename the column in the table, not only in the producer

A rename is safe in Iceberg because the field ID stays the same. Rename the table column first, then let the producer ship the new name. On Spark, the Iceberg DDL documentation gives the statement.

ALTER TABLE prod.db.events RENAME COLUMN cust_id TO customer_id;

For a period where old and new records are both on the topic, a name mapping can point both names at the same field ID. The spec allows several names per field ID in schema.name-mapping.default, and the sink’s field lookup uses that mapping when the table has one. A column listed in the mapping matches only the names given for it, so include the current name as well.

ALTER TABLE prod.db.events SET TBLPROPERTIES (
  'schema.name-mapping.default' = '[
    {"field-id": 3, "names": ["customer_id", "cust_id"]}
  ]'
);

Check the field ID in the table’s current schema before writing the mapping. The mapping shown covers only one column, and a real one should list every column.

Widen the type before the producer does

For int to long, float to double or a larger decimal precision, widen the table column first.

ALTER TABLE prod.db.events ALTER COLUMN amount TYPE double;

For any change outside the promotion list, add a new column under a new name and have the producer write both for a period.

Backfill without duplicates

The gap is in the one table copy that every engine reads, so backfill it rather than accept it. The records written during the gap are still on the topic only while the topic’s retention keeps them.

  1. Check retention. The topic must still hold the records from the producer’s release onward.
  2. Replay the topic into a staging table. A separate connector under a new name keeps the replay away from the live table.
  3. Add the source position to each row. The sink’s KafkaMetadata transform, which the sink documentation marks as experimental, adds the topic, partition, offset and timestamp.
  4. Merge the staging rows into the live table with MERGE INTO, limited to the window between the producer’s release and the fix.
  5. Compact the files the backfill wrote, using the compaction guide.

To enable the transform, set it on the staging connector.

transforms=kafkaMeta
transforms.kafkaMeta.type=org.apache.iceberg.connect.transforms.KafkaMetadataTransform

With the default settings it adds four top-level columns.

_kafka_metadata_topic
_kafka_metadata_partition
_kafka_metadata_offset
_kafka_metadata_timestamp

A merge on topic, partition and offset updates the missing values in place and inserts nothing twice. Without those columns, merge on the table’s identifier columns, which the sink takes from this setting.

iceberg.tables.default-id-columns=order_id

The merge is a commit like any other, so run it when the live writer is quiet, or expect the retries described on the commit failures page.

Where Kpow and Flex fit

This page is published by Factor House, which makes Kpow, its Kafka management product, and Flex, its Flink job management and monitoring product, so weigh this section with that in mind.

Kafka Connect. Kpow’s Kafka Connect management documentation describes restarting individual connect tasks and viewing the stacktraces of tasks in an ERROR state. Kpow does not read the Iceberg table, so it cannot show that a column stopped filling in. That check stays with the table and its metadata tables.

Schema registry. Kpow’s schema management documentation describes changing a subject’s compatibility with an Update Compatibility button, making revisions to a subject’s schema, and viewing and permanently deleting soft-deleted schemas in Confluent’s registry. It does not change the Iceberg table or the sink settings.

Flink jobs. Flex’s job inspection documentation describes an Events tab that logs the job’s lifecycle, including the job entering the RESTARTING state, which is how a Flink writer failing on a new schema first shows up. Flex does not show the Iceberg table schema or the schema updates the sink makes. The comparison of tools for self-managed Apache Flink sets it beside the Flink web UI and the REST API.

FAQ

Why does a new producer field not appear in the Iceberg table?

The Kafka Connect Iceberg sink has iceberg.tables.evolve-schema-enabled set to false by default, and with it off the sink skips the values of any field that has no matching column. Turn evolution on, or add the column to the table, then backfill the gap from the topic.

Does Iceberg support renaming a column?

Yes. Iceberg tracks columns by field ID, and a rename changes the name but keeps the ID, so existing data stays attached. The writers are the limit, because the Kafka Connect sink and the Flink Dynamic Sink both match fields by name and treat a renamed field as a new column.

Which Schema Registry compatibility mode is safe for an Iceberg sink?

No mode protects the table on its own. BACKWARD, the Confluent default, accepts deleted fields and new optional fields, both of which change what lands in the table. Pair the mode with sink settings, and make renames and type changes in the table before the producer ships them.

Can an Iceberg column change from int to string?

No. The Iceberg spec allows only widening promotions such as int to long and float to double. For any other change, add a new column, write both for a period, and move readers across.

Related reading