{"record":{"id":"08600c4f0cd5b45c","repo":"apache/iceberg","slug":"unknown-version-08600c","errorCode":null,"errorMessage":"Unknown version: ","messagePattern":"Unknown version: ","errorType":"exception","errorClass":"IOException","httpStatus":null,"severity":"error","filePath":"flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/source/enumerator/IcebergEnumeratorStateSerializer.java","lineNumber":69,"sourceCode":"  @Override\n  public int getVersion() {\n    return VERSION;\n  }\n\n  @Override\n  public byte[] serialize(IcebergEnumeratorState enumState) throws IOException {\n    return serializeV2(enumState);\n  }\n\n  @Override\n  public IcebergEnumeratorState deserialize(int version, byte[] serialized) throws IOException {\n    switch (version) {\n      case 1:\n        return deserializeV1(serialized);\n      case 2:\n        return deserializeV2(serialized);\n      default:\n        throw new IOException(\"Unknown version: \" + version);\n    }\n  }\n\n  @VisibleForTesting\n  byte[] serializeV1(IcebergEnumeratorState enumState) throws IOException {\n    DataOutputSerializer out = SERIALIZER_CACHE.get();\n    serializeEnumeratorPosition(out, enumState.lastEnumeratedPosition(), positionSerializer);\n    serializePendingSplits(out, enumState.pendingSplits(), splitSerializer);\n    byte[] result = out.getCopyOfBuffer();\n    out.clear();\n    return result;\n  }\n\n  @VisibleForTesting\n  IcebergEnumeratorState deserializeV1(byte[] serialized) throws IOException {\n    DataInputDeserializer in = new DataInputDeserializer(serialized);\n    IcebergEnumeratorPosition enumeratorPosition =\n        deserializeEnumeratorPosition(in, positionSerializer);","sourceCodeStart":51,"sourceCodeEnd":87,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/source/enumerator/IcebergEnumeratorStateSerializer.java#L51-L87","documentation":"IcebergEnumeratorStateSerializer.deserialize reads a version byte from the serialized enumerator state and dispatches to a known deserializer. This IOException is thrown when the version byte is not 1 or 2, meaning the state was written by an incompatible (newer or corrupt) serializer.","triggerScenarios":"Restoring a Flink source enumerator from checkpoint/savepoint state whose version byte is greater than 2, or reading corrupted/random bytes as serializer input.","commonSituations":"Upgrading the Iceberg Flink connector after job state was written by a newer connector version; downgrading Flink/Iceberg versions so an old job cannot read new savepoints; corrupted checkpoint files.","solutions":["Use a connector version equal to or newer than the one that wrote the checkpoint/savepoint state","Verify the state byte[] actually came from IcebergEnumeratorStateSerializer.serialize (not another operator's state)","Restore from an older compatible savepoint/checkpoint","Inspect the first byte(s) of the serialized payload to confirm the version value"],"exampleFix":"// before: restoring new-version state with old connector\n// flink-connector-iceberg 1.7 reading state serialized with version 3\n// after: align versions\n// implementation('org.apache.iceberg:iceberg-flink-runtime-1.19:1.10.0') // >= writing version","handlingStrategy":"try-catch","validationCode":"// before restore, verify state provenance\nif (stateBytes == null || stateBytes.length == 0) {\n  throw new IllegalStateException(\"Empty enumerator state; cannot restore\");\n}\n// version byte is read first by the serializer; ensure it is 1 or 2","typeGuard":null,"tryCatchPattern":"try {\n  enumeratorState = IcebergEnumeratorStateSerializer.deserialize(stateVersion, stateBytes);\n} catch (IOException e) {\n  // fall back to a fresh enumeration or fail the job with a clear message\n  throw new IllegalStateException(\"Incompatible enumerator state; recreate savepoint with matching connector version\", e);\n}","preventionTips":["Pin connector versions so the restoring version is >= the version that wrote the savepoint","Test savepoint compatibility in CI before upgrading connectors","Never hand-edit or truncate checkpoint/savepoint state bytes"],"tags":["flink","serialization","checkpoint-restore","version-incompatibility"],"backgroundTag":"invalid-enum-value","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}