{"record":{"id":"ef4b0085ccc31b51","repo":"apache/iceberg","slug":"unrecognized-version-or-corrupt-state-version-ef4b00","errorCode":null,"errorMessage":"Unrecognized version or corrupt state: \" + version","messagePattern":"Unrecognized version or corrupt state: \" \\+ version","errorType":"exception","errorClass":"IOException","httpStatus":null,"severity":"error","filePath":"flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicCommittableSerializer.java","lineNumber":70,"sourceCode":"    view.writeInt(numManifests);\n    for (int i = 0; i < numManifests; i++) {\n      byte[] manifest = committable.manifests()[i];\n      view.writeInt(manifest.length);\n      view.write(manifest);\n    }\n\n    return out.toByteArray();\n  }\n\n  @Override\n  public DynamicCommittable deserialize(int version, byte[] serialized) throws IOException {\n    if (version == VERSION_1) {\n      return deserializeV1(serialized);\n    } else if (version == VERSION_2) {\n      return deserializeV2(serialized);\n    }\n\n    throw new IOException(\"Unrecognized version or corrupt state: \" + version);\n  }\n\n  private DynamicCommittable deserializeV1(byte[] serialized) throws IOException {\n    DataInputDeserializer view = new DataInputDeserializer(serialized);\n    WriteTarget key = WriteTarget.deserializeFrom(view);\n    String jobId = view.readUTF();\n    String operatorId = view.readUTF();\n    long checkpointId = view.readLong();\n    int manifestLen = view.readInt();\n    byte[] manifestBuf = new byte[manifestLen];\n    view.read(manifestBuf);\n    return new DynamicCommittable(\n        new TableKey(key.tableName(), key.branch()),\n        new byte[][] {manifestBuf},\n        jobId,\n        operatorId,\n        checkpointId);\n  }","sourceCodeStart":52,"sourceCodeEnd":88,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicCommittableSerializer.java#L52-L88","documentation":"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.","triggerScenarios":"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.","commonSituations":"Rolling upgrades where a newer client wrote state and an older job tries to restore it, or state backend corruption.","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"],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"// pin the Iceberg Flink version to match the checkpoint producer\nObjects.equals(icebergRuntimeVersion, checkpointWriterVersion);","typeGuard":null,"tryCatchPattern":"try {\n  DynamicCommittable c = serializer.deserialize(version, bytes);\n} catch (IOException e) {\n  throw new JobRecoveryFailure(\"Incompatible committable state; restart without this savepoint\", e);\n}","preventionTips":["Keep Iceberg runtime versions uniform across job upgrades","Take fresh savepoints after each version bump","Never manually edit serialized committable state"],"tags":["flink","serialization","checkpoint"],"backgroundTag":"invalid-argument-format","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"}