apache/iceberg · error · IOException
Unrecognized version or corrupt state: " + version
Error message
Unrecognized version or corrupt state: " + version
What it means
DynamicCommittableSerializer.deserialize reads a version byte and dispatches to deserializeV1 or deserializeV2. Any other version byte indicates the payload was produced by an incompatible (newer or unknown) serializer version or the bytes are corrupt, so an IOException is thrown.
Solutions
- Run the same or newer Iceberg/Flink version as the one that wrote the checkpoint
- Discard the incompatible savepoint and restart with a fresh checkpoint
- Verify state backend integrity if corruption is suspected
Defensive patterns
Strategy: try-catch
Validate before calling
// pin the Iceberg Flink version to match the checkpoint producer Objects.equals(icebergRuntimeVersion, checkpointWriterVersion);
Try / catch
try {
DynamicCommittable c = serializer.deserialize(version, bytes);
} catch (IOException e) {
throw new JobRecoveryFailure("Incompatible committable state; restart without this savepoint", e);
} Prevention
- Keep Iceberg runtime versions uniform across job upgrades
- Take fresh savepoints after each version bump
- Never manually edit serialized committable state
When it happens
Trigger: Restoring a Flink checkpoint/savepoint produced by a newer Iceberg version, or reading corrupt committable bytes from state where the version byte is not 1 or 2.
Common situations: Rolling upgrades where a newer client wrote state and an older job tries to restore it, or state backend corruption.
Understand the failure class
Background: "Invalid ... format", "must be in format X", "does not look like a ..." — invalid argument format errors across CLI tools and libraries — this error's family across 17 libraries.
Related errors
- Could not deserialize the WriteResult object
- Failed to deserialize IcebergSourceSplit. Encountered…
- Failed to deserialize IcebergSourceSplit. Encountered…
- Failed to deserialize the split.
- Unknown serialize version:
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/ef4b0085ccc31b51.
Report an issue: GitHub.
Appendix: source
Thrown at flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicCommittableSerializer.java:70
view.writeInt(numManifests);
for (int i = 0; i < numManifests; i++) {
byte[] manifest = committable.manifests()[i];
view.writeInt(manifest.length);
view.write(manifest);
}
return out.toByteArray();
}
@Override
public DynamicCommittable deserialize(int version, byte[] serialized) throws IOException {
if (version == VERSION_1) {
return deserializeV1(serialized);
} else if (version == VERSION_2) {
return deserializeV2(serialized);
}
throw new IOException("Unrecognized version or corrupt state: " + version);
}
private DynamicCommittable deserializeV1(byte[] serialized) throws IOException {
DataInputDeserializer view = new DataInputDeserializer(serialized);
WriteTarget key = WriteTarget.deserializeFrom(view);
String jobId = view.readUTF();
String operatorId = view.readUTF();
long checkpointId = view.readLong();
int manifestLen = view.readInt();
byte[] manifestBuf = new byte[manifestLen];
view.read(manifestBuf);
return new DynamicCommittable(
new TableKey(key.tableName(), key.branch()),
new byte[][] {manifestBuf},
jobId,
operatorId,
checkpointId);
}View on GitHub (pinned to 86d9c8fc54)