{"record":{"id":"83fc6527ee3ff3fe","repo":"apache/iceberg","slug":"failed-to-restore-committer-state-this-can-happen-83fc65","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/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java","lineNumber":173,"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":155,"sourceCodeEnd":191,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java#L155-L191","documentation":"IcebergFilesCommitter restores its pending checkpoints and flink jobId state in initializeState. If isRestored() is true but the jobId state is empty, the operator state could not be restored (changed operator uid with allowNonRestoredState) and the committer starts fresh, logging this detailed warning - previously pending files may be orphaned.","triggerScenarios":"initializeState during job restore where context.isRestored() but jobIdState.get() is null or empty: operator UID changed (topology change, missing uidPrefix) combined with Flink's allowNonRestoredState=true.","commonSituations":"Restarting after modifying the job graph without a stable uidPrefix; restoring a savepoint with allowNonRestoredState that silently dropped this operator's state; re-running the job from scratch against existing checkpoint files.","solutions":["Set an explicit stable uidPrefix via FlinkSink.Builder#uidPrefix(...) so the committer operator UID survives topology changes","Restore with allowNonRestoredState=false to fail fast instead of silently losing committer state","Run DeleteOrphanFiles afterward to clean files written by the lost pending commits","Re-run the job from the last successful checkpoint with the original topology"],"exampleFix":"// before\nFlinkSink.forRowData(input).forRowData(...)...append();\n// after\nFlinkSink.forRowData(input)\n    .uidPrefix(\"my-iceberg-sink\")\n    ...append();","handlingStrategy":"validation","validationCode":"// ensure a stable operator uid before appending the sink\nPreconditions.checkState(uidPrefix != null, \"Set FlinkSink.Builder#uidPrefix for stable committer state\");\n// and restore without dropping state\nenv.getCheckpointConfig().setTolerableCheckpointFailureNumber(...);\n// run restore with: flink run -s <savepoint> (WITHOUT --allowNonRestoredState)","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Always set uidPrefix explicitly on FlinkSink builders","Never use --allowNonRestoredState for jobs with Iceberg sinks","Run DeleteOrphanFiles after any restore that lost committer state"],"tags":["flink","iceberg","state-restore","checkpoint"],"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"}