{"record":{"id":"d0331f7a60e17824","repo":"apache/iceberg","slug":"failed-to-restore-committer-state-this-can-happen-d0331f","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.2/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.2/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java#L155-L191","documentation":"IcebergFilesCommitter.initializeState warns that no committer state (pending data files / job id) could be restored even though the operator is marked as restored. This typically happens when the operator uid changed and the restore ran with allowNonRestoredState, so the state was silently dropped; committed data may have been written but never committed to Iceberg.","triggerScenarios":"Resuming from a savepoint/checkpoint where the committer operator's uid differs (topology changed, uidPrefix not set) while --allowNonRestoredState is enabled, causing jobIdState.getListState to return an empty/null iterable.","commonSituations":"Changing FlinkSink topology (adding/removing operators, parallelism changes with different uids) and restarting from an old savepoint; upgrading the job without a stable uidPrefix; careless savepoint with -n/--allowNonRestoredState.","solutions":["Set an explicit uidPrefix via FlinkSink.Builder#uidPrefix so the committer operator uid is stable across topology changes.","Restore from a savepoint/checkpoint taken with the same operator uid; avoid --allowNonRestoredState for production restores.","If data files were written but uncommitted, run removeOrphanFiles to clean them and restart with correct uids from a fresh snapshot."],"exampleFix":"// before\nFlinkSink.forRowData(input).tableLoader(loader).append();\n// after\nFlinkSink.forRowData(input).tableLoader(loader).uidPrefix(\"my-iceberg-sink\").append();","handlingStrategy":"validation","validationCode":"// verify the savepoint contains committer state before resuming\n// and always set a stable operator uid\nFlinkSink.forRowData(input).uidPrefix(\"iceberg-committer\")...","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Always set uidPrefix on the sink builder","Never use --allowNonRestoredState for production restores","After topology changes, take a new savepoint before resuming","Check savepoint metadata for the committer operator uid before restart"],"tags":["flink","savepoint","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"}