{"record":{"id":"9f3ca1e1e03abcae","repo":"apache/iceberg","slug":"unrecognized-version-or-corrupt-state-9f3ca1","errorCode":null,"errorMessage":"Unrecognized version or corrupt state: ","messagePattern":"Unrecognized version or corrupt state: ","errorType":"exception","errorClass":"IOException","httpStatus":null,"severity":"error","filePath":"flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/MonitorSource.java","lineNumber":202,"sourceCode":"          iterator.getClass());\n\n      TableChangeIterator tableChangeIterator = (TableChangeIterator) iterator;\n      DataOutputSerializer out = new DataOutputSerializer(8);\n      long toStore =\n          tableChangeIterator.lastSnapshotId != null ? tableChangeIterator.lastSnapshotId : -1L;\n      out.writeLong(toStore);\n      return out.getCopyOfBuffer();\n    }\n\n    @Override\n    public TableChangeIterator deserialize(int version, byte[] serialized) throws IOException {\n      if (version == CURRENT_VERSION) {\n        DataInputDeserializer in = new DataInputDeserializer(serialized);\n        long fromStore = in.readLong();\n        return new TableChangeIterator(\n            tableLoader, fromStore != -1 ? fromStore : null, maxReadBack);\n      } else {\n        throw new IOException(\"Unrecognized version or corrupt state: \" + version);\n      }\n    }\n  }\n}\n","sourceCodeStart":184,"sourceCodeEnd":207,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/MonitorSource.java#L184-L207","documentation":"MonitorSource's TableChangeIterator state serializer writes a version byte/integer into the serialized operator state. On deserialize, if the stored version does not equal CURRENT_VERSION, an IOException is thrown because the state layout cannot be interpreted. This guards against restoring state written by an incompatible serializer version or corrupted state.","triggerScenarios":"Restoring a Flink savepoint/checkpoint that contains MonitorSource state serialized by a different Iceberg version (different CURRENT_VERSION), or state storage corruption producing a wrong version field.","commonSituations":"Upgrading Iceberg/Flink between jobs, restoring old checkpoints into a new deployment, or manually edited/truncated checkpoint files.","solutions":["Restore from a checkpoint/savepoint created with the same Iceberg version as the running job.","If upgrade is required, let the old job finish or start the maintenance source fresh (allowing non-restored state) instead of restoring incompatible state.","Verify checkpoint files are not truncated or corrupted; re-run from an earlier valid checkpoint."],"exampleFix":null,"handlingStrategy":"fallback","validationCode":"// Before restoring, check state was produced by a compatible version:\n// compare the Iceberg version that wrote the savepoint with the running version.","typeGuard":null,"tryCatchPattern":"try {\n  iterator = TableChangeIterator.deserialize(tableLoader, serialized, maxReadBack);\n} catch (IOException e) {\n  LOG.warn(\"Incompatible MonitorSource state; reinitializing from scratch\", e);\n  iterator = TableChangeIterator.initial(tableLoader, maxReadBack);\n}","preventionTips":["Restore savepoints only with the same Iceberg version that produced them.","Avoid manual edits to checkpoint metadata; verify checkpoint integrity."],"tags":["flink","serialization","state-restore"],"backgroundTag":"schema-validation-failed","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"}