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

  1. Use the Iceberg version matching the savepoint's serialization version to restore
  2. Rebuild job state without restore (start fresh; committed data remains in the table)
  3. 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

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


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/91b59ee5428536d0. Report an issue: GitHub.