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

  1. Run the same or newer Iceberg/Flink version as the one that wrote the checkpoint
  2. Discard the incompatible savepoint and restart with a fresh checkpoint
  3. 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

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


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)