apache/iceberg · error · IOException
Unrecognized version or corrupt state: ${version}
Error message
Unrecognized version or corrupt state: ${version} What it means
DynamicWriteResultSerializer.deserialize reads a version byte and delegates to the underlying WriteResult serializer only for versions it knows; any other version throws this IOException. Like the committable serializer, it protects against decoding write-result state produced by incompatible Iceberg versions.
Source
Thrown at flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicWriteResultSerializer.java:62
view.writeInt(writeResult.specId());
byte[] result = WRITE_RESULT_SERIALIZER.serialize(writeResult.writeResult());
view.write(result);
return out.toByteArray();
}
@Override
public DynamicWriteResult deserialize(int version, byte[] serialized) throws IOException {
if (version == 1) {
DataInputDeserializer view = new DataInputDeserializer(serialized);
TableKey key = TableKey.deserializeFrom(view);
int specId = view.readInt();
byte[] resultBuf = new byte[view.available()];
view.read(resultBuf);
WriteResult writeResult = WRITE_RESULT_SERIALIZER.deserialize(version, resultBuf);
return new DynamicWriteResult(key, specId, writeResult);
}
throw new IOException("Unrecognized version or corrupt state: " + version);
}
}
View on GitHub (pinned to 86d9c8fc54)
Solutions
- Use the Iceberg version matching the savepoint's serialization version to restore
- Rebuild job state without restore (start fresh; committed data remains in the table)
- Check for classpath mixing of Iceberg Flink artifacts and align all nodes to one version
Defensive patterns
Strategy: try-catch
Validate before calling
// Align iceberg-flink-runtime versions between the job that wrote the state and the restoring job.
Try / catch
try {
result = serializer.deserialize(version, bytes);
} catch (IOException e) {
// unknown version byte: restore with the matching Iceberg version or start without state
} Prevention
- Keep JobManager and TaskManager Iceberg artifacts identical
- Restore savepoints with the Iceberg version that wrote them
- Exercise cross-version restore tests before upgrades
When it happens
Trigger: Restoring a Flink job whose operator state contains DynamicWriteResult bytes written with an unknown version byte (different/newer Iceberg sink version, or corrupted state). Called transitively via copy() when Flink duplicates state or manually via testUnsupportedVersion.
Common situations: Upgrading the Iceberg Flink sink between incompatible serializer versions and resuming from an old savepoint; mixed Iceberg jar versions across JobManager/TaskManagers; corrupted checkpoint bytes.
Related errors
- Unknown serialize version: ${version}
- Unrecognized version or corrupt state: ${version}
- Failed to deserialize IcebergSourceSplit. Encountered unsupp
- Unrecognized version or corrupt state: <version>
- Unrecognized version or corrupt state: ${version}
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/91b59ee5428536d0.
Report an issue: GitHub.