apache/iceberg · error · IOException
Unrecognized version or corrupt state:
Error message
Unrecognized version or corrupt state:
What it means
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.
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.
Defensive patterns
Strategy: fallback
Validate before calling
// Before restoring, check state was produced by a compatible version: // compare the Iceberg version that wrote the savepoint with the running version.
Try / catch
try {
iterator = TableChangeIterator.deserialize(tableLoader, serialized, maxReadBack);
} catch (IOException e) {
LOG.warn("Incompatible MonitorSource state; reinitializing from scratch", e);
iterator = TableChangeIterator.initial(tableLoader, maxReadBack);
} Prevention
- Restore savepoints only with the same Iceberg version that produced them.
- Avoid manual edits to checkpoint metadata; verify checkpoint integrity.
When it happens
Trigger: 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.
Common situations: Upgrading Iceberg/Flink between jobs, restoring old checkpoints into a new deployment, or manually edited/truncated checkpoint files.
Understand the failure class
Background: Schema validation failed / invalid input schema: payload rejected because its shape doesn't match the expected schema — this error's family across 28 libraries.
Related errors
- Failed to initialize serializerCache for reading data with…
- Unknown read version:
- Unknown read version
- Unknown serialize version:
- Could not deserialize the WriteResult object
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/9f3ca1e1e03abcae.
Report an issue: GitHub.
Appendix: source
Thrown at flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/MonitorSource.java:202
iterator.getClass());
TableChangeIterator tableChangeIterator = (TableChangeIterator) iterator;
DataOutputSerializer out = new DataOutputSerializer(8);
long toStore =
tableChangeIterator.lastSnapshotId != null ? tableChangeIterator.lastSnapshotId : -1L;
out.writeLong(toStore);
return out.getCopyOfBuffer();
}
@Override
public TableChangeIterator deserialize(int version, byte[] serialized) throws IOException {
if (version == CURRENT_VERSION) {
DataInputDeserializer in = new DataInputDeserializer(serialized);
long fromStore = in.readLong();
return new TableChangeIterator(
tableLoader, fromStore != -1 ? fromStore : null, maxReadBack);
} else {
throw new IOException("Unrecognized version or corrupt state: " + version);
}
}
}
}
View on GitHub (pinned to 86d9c8fc54)