{"record":{"id":"0ff2e5c8a94a8e7c","repo":"apache/iceberg","slug":"failed-to-restore-committer-state-this-can-happen-0ff2e5","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.3/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.3/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java#L155-L191","documentation":"A warning logged when the IcebergFilesCommitter operator restores from a checkpoint/savepoint but finds no persisted job ID state. This usually means the operator UID changed between jobs and Flink's allowNonRestoredState allowed the restore to proceed without the committer's state, so in-flight files from before cannot be committed.","triggerScenarios":"Restoring a Flink job whose topology changed (altering the auto-generated operator UID) with execution.checkpointing.allowed-state-restoration-types or allowNonRestoredState enabled; ListState for IcebergCommitState/JOB_ID_DESCRIPTOR is null or empty on restore.","commonSituations":"Job topology refactors (added/removed sink operators, renamed builders), missing uidPrefix() on the sink builder, upgrading Iceberg versions changing operator chains, restoring across Flink versions.","solutions":["Set an explicit stable operator uid via FlinkSink.Builder#uidPrefix(\"my-sink-uid\") so the committer operator uid survives topology changes","Restore from a checkpoint taken before the topology change, or without allowNonRestoredState so mismatches fail loudly instead of losing state","After a UID mismatch, expect orphaned data files from previous pending commits; clean them via expire/cleanup or a manual orphan-file removal"],"exampleFix":"// before\nFlinkSink.forRowData(input).append();\n// after\nFlinkSink.forRowData(input).uidPrefix(\"iceberg-committer-stable\").append();","handlingStrategy":"validation","validationCode":"// Ensure a stable uid before building the sink\nPreconditions.checkNotNull(uidPrefix, \"uidPrefix must be set for the Iceberg committer\");","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Always call uidPrefix() with a stable value on FlinkSink/IcebergSink builders","Avoid allowNonRestoredState in production restores so state mismatches fail loudly","Keep job topology changes and restore points coordinated"],"tags":["flink","checkpoint","state-restoration","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"}