{"record":{"id":"a0695984f9bd5ab6","repo":"apache/iceberg","slug":"failed-to-restore-committer-state-this-can-happen","errorCode":null,"errorMessage":"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.","messagePattern":"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\\.","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java","lineNumber":170,"sourceCode":"\n    maxContinuousEmptyCommits =\n        PropertyUtil.propertyAsInt(table.properties(), MAX_CONTINUOUS_EMPTY_COMMITS, 10);\n    Preconditions.checkArgument(\n        maxContinuousEmptyCommits > 0, MAX_CONTINUOUS_EMPTY_COMMITS + \" must be positive\");\n\n    int subTaskId = getRuntimeContext().getTaskInfo().getIndexOfThisSubtask();\n    int attemptId = getRuntimeContext().getTaskInfo().getAttemptNumber();\n    this.manifestOutputFileFactory =\n        FlinkManifestUtil.createOutputFileFactory(\n            () -> table, table.properties(), flinkJobId, operatorUniqueId, subTaskId, attemptId);\n    this.maxCommittedCheckpointId = INITIAL_CHECKPOINT_ID;\n\n    this.checkpointsState = context.getOperatorStateStore().getListState(STATE_DESCRIPTOR);\n    this.jobIdState = context.getOperatorStateStore().getListState(JOB_ID_DESCRIPTOR);\n    if (context.isRestored()) {\n      Iterable<String> jobIdIterable = jobIdState.get();\n      if (jobIdIterable == null || !jobIdIterable.iterator().hasNext()) {\n        LOG.warn(\n            \"Failed to restore committer state. This can happen when operator uid changed and Flink \"\n                + \"allowNonRestoredState is enabled. Best practice is to explicitly set the operator id \"\n                + \"via FlinkSink#Builder#uidPrefix() so that the committer operator uid is stable. \"\n                + \"Otherwise, Flink auto generate an operator uid based on job topology.\"\n                + \"With that, operator uid is subjective to change upon topology change.\");\n        return;\n      }\n\n      String restoredFlinkJobId = jobIdIterable.iterator().next();\n      Preconditions.checkState(\n          !Strings.isNullOrEmpty(restoredFlinkJobId),\n          \"Flink job id parsed from checkpoint snapshot shouldn't be null or empty\");\n\n      // Since flink's checkpoint id will start from the max-committed-checkpoint-id + 1 in the new\n      // flink job even if it's restored from a snapshot created by another different flink job, so\n      // it's safe to assign the max committed checkpoint id from restored flink job to the current\n      // flink job.\n      this.maxCommittedCheckpointId =","sourceCodeStart":152,"sourceCodeEnd":188,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java#L152-L188","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","solutions":["Set an explicit stable uid via FlinkSink.builder().uidPrefix(\"my-iceberg-committer\") and re-save a new savepoint.","Disable allowNonRestoredState so restores fail loudly instead of silently dropping committer state.","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.","After migrating, immediately trigger a new savepoint with the stable uid so future restores are safe."],"exampleFix":"// before\nFlinkSink.forRowData(input)\n    .table(table)\n    .append();\n// after\nFlinkSink.forRowData(input)\n    .table(table)\n    .uidPrefix(\"iceberg-files-committer\")\n    .append();","handlingStrategy":"validation","validationCode":"// before submitting\nboolean stableUid = sinkBuilder != null && sinkBuilderUidPrefixSet;\nif (restoringFromSavepoint && !stableUid) {\n  throw new IllegalStateException(\"Set uidPrefix before restoring Iceberg sink from a savepoint\");\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Always set uidPrefix on FlinkSink builders in production jobs.","Keep allowNonRestoredState=false so uid mismatches fail fast.","Re-save savepoints after any topology change."],"tags":["flink","checkpointing","state-restore","operator-uid"],"backgroundTag":"invalid-state-transition","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"}