Skip to content

Iceberg commits failing or conflicting: find and fix the cause

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

The alert on the Iceberg table says no new data has landed for an hour. The Flink job writing it is running, or it is restarting in a loop, and its logs show org.apache.iceberg.exceptions.CommitFailedException. On a Kafka Connect pipeline, the Iceberg sink’s commit coordinator has stopped after a failed commit, and consumer lag on the source topics keeps climbing.

A failed Iceberg commit is a writer that has written its data files but could not make them part of the table, because another commit changed the table first and the writer ran out of retries, or because it could not confirm whether its commit was applied. The Iceberg Javadoc describes CommitFailedException as an exception raised when a commit fails because of out of date metadata.

This page is about the commit step itself. If the Flink job’s checkpoints fail or time out before any commit is attempted, start with Flink checkpoints that fail or time out. If commits succeed but the table holds thousands of tiny files, see too many small files and how to compact them. For how tables, snapshots and metadata fit together, see the complete Iceberg guide.

Rule out a commit that never ran

On Flink, a missing commit is not always a failed one. The Flink writes documentation says the Iceberg commit happens after a successful Flink checkpoint, in the notifyCheckpointComplete callback. It also says Iceberg commits can fail while Flink checkpoints succeed, and that if the callback is not triggered, no Iceberg commit is attempted at all. The committer metric below, under the IcebergFilesCommitter metric group, is the one the documentation calls ideal for alerting on failed or missing commits.

elapsedSecondsSinceLastSuccessfulCommit

The documentation’s example is a 5 minute checkpoint interval with an alert when this metric exceeds 60 minutes. If the metric climbs and the job’s checkpoints are also failing, the commit was never reached, and the Flink checkpoint troubleshooting guide is the place to start. If checkpoints complete and the metric still climbs, the commit itself is failing, and the rest of this page applies.

Two per-job signals show whether the Flink job is healthy enough to reach a checkpoint at all, backpressure and checkpoint duration. The Flink back pressure monitoring documentation covers the first, and the checkpoint monitoring documentation defines the end to end duration of each checkpoint.

Which commit failure you are looking at

Two exceptions look alike in a stack trace and need different fixes. A CommitFailedException is a conflict that retries did not resolve. A CommitStateUnknownException, which the Javadoc describes as a failure to confirm either affirmatively or negatively that a commit was applied, is a lost connection during the commit. The configuration documentation says commit.status-check.num-retries, 3 by default, is the number of times to check whether a commit succeeded after a connection is lost before failing due to an unknown commit state.

commit.status-check.num-retries=3
commit.status-check.min-wait-ms=1000
commit.status-check.max-wait-ms=60000
commit.status-check.total-timeout-ms=1800000

An unknown state is a catalog or network problem, so raising the conflict retries does not help it. The same Javadoc says the client cannot take any further action without possibly corrupting the table, so check whether the writer’s snapshot landed before restarting anything, using the snapshots query below.

F1 A conflicting commit and an unknown commit state, side by side
Conflict Unknown state
Exception CommitFailedException CommitStateUnknownException
What it means The commit failed because the writer's table metadata was out of date. The writer could not confirm whether the commit was applied or not.
Usual cause Another writer or a maintenance job committed first, and retries ran out. The connection to the catalog was lost during the commit.
Properties that govern it The commit.retry.* properties The commit.status-check.* properties
First check The snapshots table, for commits by other jobs around the failure time. The snapshots table, for whether the writer's snapshot landed.
Source: Iceberg table configuration documentation and the Iceberg Javadoc

Find the writer or job that conflicts

List every commit around the failure

The snapshots metadata table has one row per commit, with its time, operation and summary. The Iceberg Spark queries documentation lists it among the metadata tables that ship with Iceberg, so the query below needs no extra tooling. Query the window around the first failure from any engine that reads Iceberg metadata tables.

SELECT committed_at, snapshot_id, operation,
       summary['flink.job-id'] AS flink_job
FROM prod.db.events.snapshots
WHERE committed_at > TIMESTAMP '2026-10-03 09:00:00'
ORDER BY committed_at

The table specification defines the operation values. append means only data files were added. replace means files were added and removed without changing table data, for example compaction. overwrite and delete come from logical overwrites and deletes. The Flink writes documentation says Flink streaming write jobs identify their snapshots with the flink.job-id property in the summary. A replace or overwrite commit that lands between the writer’s last success and its failure shows the kind of commit that won the race. Two different flink.job-id values whose commits interleave in time mean two Flink jobs are writing the table.

The specification requires only operation in a snapshot’s summary and lists the other summary fields as optional, so an empty flink_job means the snapshot came from something other than a Flink streaming write job, or from a writer that does not populate that key. It also helps to read the latest snapshot time rather than the table’s last updated time. The specification says each table metadata file updates last-updated-ms, and it describes schema updates and partition spec changes as commits that must validate against the current table version. A schema change therefore moves the table’s last updated time without adding data, while the newest row in the snapshots table is the last data commit.

SELECT max(committed_at) AS last_data_commit
FROM prod.db.events.snapshots

Read the commit metrics

The Flink writes documentation lists committer metrics beside the alerting metric. lastCommitDurationMs is the duration of the Iceberg table commit, and committedDataFilesCount is the number of data files committed.

lastCommitDurationMs
committedDataFilesCount

A commit duration that keeps growing is worth reading beside the snapshots table, because commits by other jobs in the same window mean the writer is spending that time on retries. A file count per commit that rises with writer parallelism points at the small files problem, which also makes each compaction rewrite more files and widens its window for conflicts.

On Kafka Connect, check the coordinator

The Apache Iceberg Kafka Connect sink runs a coordinator that commits on behalf of all tasks. Its configuration says iceberg.control.commit.max-consecutive-failures is the maximum number of consecutive commit failures before the coordinator terminates, with a default of 1, and that iceberg.control.commit.timeout-ms defaults to 30,000 ms.

iceberg.control.commit.interval-ms=300000
iceberg.control.commit.timeout-ms=30000
iceberg.control.commit.max-consecutive-failures=1

With the default of 1, a single failed commit ends the coordinator, which matches the symptom of commits stopping while source lag climbs. The Kafka Connect REST status endpoint returns the state of the connector and every task, with error information, which the guide to diagnosing a failed Kafka Connect connector walks through. Read the status first, then look for the competing writer in the snapshots table.

GET /connectors/{name}/status

The Flink Kafka connector documentation says the Kafka source commits its consuming offsets when checkpoints are completed, and that it does not rely on those offsets for fault tolerance, committing them only to expose progress for monitoring. Lag on the job’s consumer group therefore moves with completed checkpoints, not with records read.

commit.offsets.on.checkpoint

Lag that grows between checkpoints and drops as each one completes is the expected pattern. Lag that stops dropping means offsets are no longer being committed, which points back to checkpoints that are not completing. Lag that drops at each checkpoint while elapsedSecondsSinceLastSuccessfulCommit climbs means checkpoints complete and the Iceberg commit is the step that is failing.

How an Iceberg commit works

Iceberg has no lock on a table while a writer prepares its change. Every writer works on its own copy of the metadata and only finds out about a conflict when it tries to commit.

One atomic swap, retried on conflict

The Iceberg reliability documentation says commits replace the path of the current table metadata file using an atomic operation, and that this is the basis for serializable isolation. The table specification describes the same model as optimistic concurrency. Writers create table metadata files assuming the current version will not change before they commit, then commit by swapping the table’s metadata file pointer from the base version to the new version. If two writers start from the same version, only one swap can succeed, and the other must retry its update on the new current version.

Retries are cheap for appends and checked for rewrites. The reliability documentation says appends usually create a new manifest file for the appended data files, which can be added to the table without rewriting the manifest on every attempt, so a streaming append that loses the race normally retries and succeeds. A rewrite cannot reuse its work the same way. The documentation says commits are structured as assumptions and actions, and after a conflict the writer checks that its assumptions still hold. Its example is a compaction that rewrites two files into one. That commit is safe only while the table still contains both source files, and if a conflicting commit deleted either of them, the operation must fail.

Retries run out on a schedule set by the table

The retry budget is a set of table properties, not a Flink or Kafka Connect setting. The table configuration documentation lists these defaults: 4 retries, a wait between 100 ms and 1 minute, and a total retry timeout of 30 minutes.

commit.retry.num-retries=4
commit.retry.min-wait-ms=100
commit.retry.max-wait-ms=60000
commit.retry.total-timeout-ms=1800000

A writer that still loses the swap after those retries, or after the total timeout, fails the commit with CommitFailedException.

More than one process commits to a streaming table

A streaming table usually has several committers: the streaming writer, a compaction job, a snapshot expiry job, and sometimes a backfill or a Spark MERGE. Each can run in a different engine, and each commits through the same metadata pointer, so any of them can be the commit that wins the swap.

Fix it, one change at a time

Match the evidence from the snapshots table to the fix, apply one change, then watch elapsedSecondsSinceLastSuccessfulCommit or the coordinator’s commit log before the next.

  • Short bursts of competing appends point to more retries or a longer commit interval.
  • Compaction commits failing or forcing writer retries point to scheduling maintenance away from the writer.
  • Two flink.job-id values with interleaved commits, or several jobs appending, point to one committer per table.
  • A CommitStateUnknownException points to the catalog or the network, not to retries.
  • Duplicate data after a restart points to idempotent commits.
  • A Spark MERGE, UPDATE or DELETE as the competing job points to the isolation level.

Give the commit more retries

When the snapshots table shows short bursts of competing commits, a larger retry budget lets the writer wait them out. The properties are set on the table, so every writer to that table picks them up.

ALTER TABLE prod.db.events SET TBLPROPERTIES (
  'commit.retry.num-retries' = '10',
  'commit.retry.min-wait-ms' = '1000'
)

The values above are an example, not an Iceberg recommendation. The worst case wait is roughly the number of retries multiplied by commit.retry.max-wait-ms, and commit.retry.total-timeout-ms caps the total. At the default maximum wait of 60,000 ms, 10 retries could wait up to 600 s, which is inside the 30 minute cap. A longer wait keeps the writer inside the commit, so watch lastCommitDurationMs after the change.

Commit less often

Each commit is one more attempt at the metadata swap, and the Flink writes documentation ties the commit to the checkpoint, so the checkpoint interval is the commit interval, and the Kafka Connect sink commits on iceberg.control.commit.interval-ms, 5 minutes by default.

execution.checkpointing.interval: 5min

At one commit every minute a table gets 1,440 commits a day, and at one every five minutes it gets 288. Fewer, larger commits also produce fewer small files, which is the subject of the small files and compaction guide.

Run compaction on partitions the writer has left

Compaction is the maintenance job the reliability documentation uses as its example of a commit that must fail after a conflict. An append by the streaming writer does not remove files, so it does not by itself break a compaction’s assumptions, while a writer that also deletes or overwrites files can. When the snapshots table shows compaction commits failing or forcing writer retries, run compaction on partitions the writer is no longer appending to. The rewrite_data_files procedure in the Spark procedures documentation accepts a where filter that selects the files to rewrite.

CALL prod.system.rewrite_data_files(
  table => 'db.events',
  where => 'event_date < current_date()',
  options => map('partial-progress.enabled', 'true')
)

The same documentation says partial-progress.enabled commits groups of files before the entire rewrite completes, with partial-progress.max-commits defaulting to 10. A smaller rewrite that commits in parts has less to lose when one part conflicts.

Snapshot expiry also commits to the table, and it can break the Flink writer in a different way. The Flink writes documentation says Flink streaming write jobs keep the last committed checkpoint ID in the snapshot summary, so the last snapshot created by the Flink job must be kept. The guide to expiring snapshots and orphan files covers the retention settings.

The Flink TableMaintenance documentation describes a lock that ensures only one maintenance operation runs per table at a time. That lock separates maintenance tasks from each other, not maintenance from the writer, so it does not replace scheduling. The Flink writes documentation also describes post-commit maintenance on the IcebergSink, and labels the SinkV2 based IcebergSink as an experimental feature to use with caution.

Give each table one committer

Two streaming jobs appending to the same table compete for every commit. The simplest structural fix is one committing writer per table. On Kafka Connect, the Apache Iceberg Kafka Connect sink, the open source and self-managed way to write Kafka topics into Iceberg, lists commit coordination for centralized Iceberg commits and multi-table fan-out among its features, so one connector can commit for all of its tasks and for several tables.

On Flink, the Dynamic Iceberg Sink adds its own source of conflicts. The Flink writes documentation says that in the forward path, schema changes are applied immediately, which can cause many conflicting commits to the Iceberg catalog and temporarily delay data processing. Its advice is to update the schema externally before publishing records with the new schema, or to plan for a temporary drop in throughput when a new schema arrives.

Writing to a separate branch separates the data but not the commit. The specification says branches are updated by committing a new snapshot through the same conflict resolution and retry procedures. The branch is set with the sink’s toBranch option on Flink and with a table-specific property on Kafka Connect.

FlinkSink.forRowData(input).toBranch("audit-branch")
iceberg.table.<table-name>.commit-branch=audit-branch

Keep retried commits idempotent

A writer that restarts after a failed commit must not add the same data twice. Both writers handle this, as long as their state survives.

On Flink, the writer records the last committed checkpoint ID in the snapshot summary. The Flink savepoints documentation highly recommends assigning operator IDs with uid, so keep the sink’s IDs stable across deployments. The Iceberg Flink writes documentation sets them with uidPrefix on FlinkSink, and says to use uidSuffix instead when using IcebergSink.

FlinkSink: uidPrefix
IcebergSink: uidSuffix

On Kafka Connect, the Iceberg sink documentation says the sink relies on KIP-447 for exactly-once semantics, which requires Kafka 2.5 or later. iceberg.coordinator.transactional.prefix sets the prefix of the transactional ID used by the coordinator’s producer.

iceberg.coordinator.transactional.prefix=iceberg-events-

To confirm the fix, run the snapshots query after the restart. Snapshots from the same flink.job-id resume, and elapsedSecondsSinceLastSuccessfulCommit falls back to its normal range.

Check the isolation level of Spark row-level commands

If the conflicting job is a Spark MERGE, UPDATE or DELETE, the table’s isolation level for that command applies. The configuration documentation lists three properties, each serializable by default with snapshot as the alternative.

write.merge.isolation-level=serializable
write.update.isolation-level=serializable
write.delete.isolation-level=serializable

Change one only after confirming the command’s correctness does not depend on serializable isolation. To confirm the change, run the Spark command again and check that it no longer ends in a CommitFailedException, and that elapsedSecondsSinceLastSuccessfulCommit on the Flink writer stays in its normal range.

Where Flex fits

This page is published by Factor House, which makes Flex, its Flink job management and monitoring product, so weigh this section with that in mind. Flex’s job inspection documentation describes a Checkpoints tab with counts by status, including failed checkpoints, average checkpoint size and duration, and a history of every checkpoint attempt, plus an Events tab that logs the job’s lifecycle. That separates a job whose checkpoints fail from one whose checkpoints complete while its commits do not. It does not show the Iceberg commit result or elapsedSecondsSinceLastSuccessfulCommit, so those stay in the committer metrics, the Iceberg metadata tables and the writer’s logs.

On a Kafka Connect pipeline, Kpow, the Factor House Kafka management product, has Kafka Connect management documentation that describes restarting individual connect tasks and viewing the stacktraces of tasks in an ERROR state.

Flex is one of several ways to watch a Flink job, and the comparison of tools for self-managed Apache Flink sets it beside the Flink web UI, the REST API and metrics stacks. The complete Flink guide covers the rest of the job side.

FAQ

What causes CommitFailedException in Iceberg?

The writer’s table metadata was out of date when it tried to commit, because another writer or a maintenance job committed first, and the writer used up its retries. The default is 4 retries within 30 minutes, set by the commit.retry table properties.

Can a Flink checkpoint succeed while the Iceberg commit fails?

Yes. The Iceberg Flink writes documentation says the commit runs after a successful checkpoint and that commits can fail while checkpoints succeed. Alert on elapsedSecondsSinceLastSuccessfulCommit, not on checkpoint failures alone.

Does compaction conflict with a streaming writer?

It can. A compaction commit is valid only while every file it rewrote is still in the table, so a writer that deletes or overwrites those files makes it fail. Running compaction on partitions the writer no longer touches, with partial progress enabled, keeps the overlap small.

Is it safe to raise commit.retry.num-retries?

Yes, within limits. More retries let a writer outlast brief contention, but if another job commits constantly the writer still fails later. The number of retries multiplied by commit.retry.max-wait-ms bounds the wait, commit.retry.total-timeout-ms caps it, and the competing schedule still needs fixing.

Related reading