apache/iceberg · warning

Failed to restore committer state. This can happen when…

Error message

Failed to restore committer state. This can happen when operator uid changed and Flink allowNonRestoredState is enabled. Best practice is to explicitly set the operator id via FlinkSink#Builder#uidPrefix() so that the committer operator uid is stable. Otherwise, Flink auto generate an operator uid based on job topology.With that, operator uid is subjective to change upon topology change.

What it means

The IcebergFilesCommitter restores its pending-file state from Flink operator state on job recovery. This warning is logged when restoration succeeded as an operator but the committer's state list is empty after a restore, meaning the previously committed state was lost. It almost always indicates the committer operator's uid changed between jobs and allowNonRestoredState allowed the snapshot to restore anyway, silently dropping state.

Solutions

  1. Set an explicit stable uid via FlinkSink.builder().uidPrefix("my-iceberg-committer") and re-save a new savepoint.
  2. Disable allowNonRestoredState so restores fail loudly instead of silently dropping committer state.
  3. Restore from a savepoint taken with the same topology/uid; if impossible, accept the state loss and ensure downstream compaction/rewrite fixes any uncommitted files.
  4. After migrating, immediately trigger a new savepoint with the stable uid so future restores are safe.

Example fix

// before
FlinkSink.forRowData(input)
    .table(table)
    .append();
// after
FlinkSink.forRowData(input)
    .table(table)
    .uidPrefix("iceberg-files-committer")
    .append();
Defensive patterns

Strategy: validation

Validate before calling

// before submitting
boolean stableUid = sinkBuilder != null && sinkBuilderUidPrefixSet;
if (restoringFromSavepoint && !stableUid) {
  throw new IllegalStateException("Set uidPrefix before restoring Iceberg sink from a savepoint");
}

Prevention

When it happens

Trigger: Restoring a Flink job from a savepoint/checkpoint where the IcebergFilesCommitter operator uid differs from the one that created the snapshot, with execution.checkpointing.allowNonRestoredState=true. This happens after topology changes (adding/removing operators, changing parallelism handling) without an explicit uidPrefix.

Common situations: Upgrading the Flink job graph or the Iceberg sink version, renaming the job's sink chain, switching between FlinkSink and other sink builders, or rebuilding the topology in code — causing Flink's auto-generated hash-based operator uid to change.

Understand the failure class

Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.

Related errors


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/a0695984f9bd5ab6. Report an issue: GitHub.

Appendix: source

Thrown at flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java:170

    maxContinuousEmptyCommits =
        PropertyUtil.propertyAsInt(table.properties(), MAX_CONTINUOUS_EMPTY_COMMITS, 10);
    Preconditions.checkArgument(
        maxContinuousEmptyCommits > 0, MAX_CONTINUOUS_EMPTY_COMMITS + " must be positive");

    int subTaskId = getRuntimeContext().getTaskInfo().getIndexOfThisSubtask();
    int attemptId = getRuntimeContext().getTaskInfo().getAttemptNumber();
    this.manifestOutputFileFactory =
        FlinkManifestUtil.createOutputFileFactory(
            () -> table, table.properties(), flinkJobId, operatorUniqueId, subTaskId, attemptId);
    this.maxCommittedCheckpointId = INITIAL_CHECKPOINT_ID;

    this.checkpointsState = context.getOperatorStateStore().getListState(STATE_DESCRIPTOR);
    this.jobIdState = context.getOperatorStateStore().getListState(JOB_ID_DESCRIPTOR);
    if (context.isRestored()) {
      Iterable<String> jobIdIterable = jobIdState.get();
      if (jobIdIterable == null || !jobIdIterable.iterator().hasNext()) {
        LOG.warn(
            "Failed to restore committer state. This can happen when operator uid changed and Flink "
                + "allowNonRestoredState is enabled. Best practice is to explicitly set the operator id "
                + "via FlinkSink#Builder#uidPrefix() so that the committer operator uid is stable. "
                + "Otherwise, Flink auto generate an operator uid based on job topology."
                + "With that, operator uid is subjective to change upon topology change.");
        return;
      }

      String restoredFlinkJobId = jobIdIterable.iterator().next();
      Preconditions.checkState(
          !Strings.isNullOrEmpty(restoredFlinkJobId),
          "Flink job id parsed from checkpoint snapshot shouldn't be null or empty");

      // Since flink's checkpoint id will start from the max-committed-checkpoint-id + 1 in the new
      // flink job even if it's restored from a snapshot created by another different flink job, so
      // it's safe to assign the max committed checkpoint id from restored flink job to the current
      // flink job.
      this.maxCommittedCheckpointId =

View on GitHub (pinned to 86d9c8fc54)