apache/iceberg · error · IOException

Unrecognized version or corrupt state:

Error message

Unrecognized version or corrupt state: 

What it means

MonitorSource's TableChangeIterator state serializer writes a version byte/integer into the serialized operator state. On deserialize, if the stored version does not equal CURRENT_VERSION, an IOException is thrown because the state layout cannot be interpreted. This guards against restoring state written by an incompatible serializer version or corrupted state.

Solutions

  1. Restore from a checkpoint/savepoint created with the same Iceberg version as the running job.
  2. If upgrade is required, let the old job finish or start the maintenance source fresh (allowing non-restored state) instead of restoring incompatible state.
  3. Verify checkpoint files are not truncated or corrupted; re-run from an earlier valid checkpoint.
Defensive patterns

Strategy: fallback

Validate before calling

// Before restoring, check state was produced by a compatible version:
// compare the Iceberg version that wrote the savepoint with the running version.

Try / catch

try {
  iterator = TableChangeIterator.deserialize(tableLoader, serialized, maxReadBack);
} catch (IOException e) {
  LOG.warn("Incompatible MonitorSource state; reinitializing from scratch", e);
  iterator = TableChangeIterator.initial(tableLoader, maxReadBack);
}

Prevention

When it happens

Trigger: Restoring a Flink savepoint/checkpoint that contains MonitorSource state serialized by a different Iceberg version (different CURRENT_VERSION), or state storage corruption producing a wrong version field.

Common situations: Upgrading Iceberg/Flink between jobs, restoring old checkpoints into a new deployment, or manually edited/truncated checkpoint files.

Understand the failure class

Background: Schema validation failed / invalid input schema: payload rejected because its shape doesn't match the expected schema — this error's family across 28 libraries.

Related errors


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

Appendix: source

Thrown at flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/MonitorSource.java:202

          iterator.getClass());

      TableChangeIterator tableChangeIterator = (TableChangeIterator) iterator;
      DataOutputSerializer out = new DataOutputSerializer(8);
      long toStore =
          tableChangeIterator.lastSnapshotId != null ? tableChangeIterator.lastSnapshotId : -1L;
      out.writeLong(toStore);
      return out.getCopyOfBuffer();
    }

    @Override
    public TableChangeIterator deserialize(int version, byte[] serialized) throws IOException {
      if (version == CURRENT_VERSION) {
        DataInputDeserializer in = new DataInputDeserializer(serialized);
        long fromStore = in.readLong();
        return new TableChangeIterator(
            tableLoader, fromStore != -1 ? fromStore : null, maxReadBack);
      } else {
        throw new IOException("Unrecognized version or corrupt state: " + version);
      }
    }
  }
}

View on GitHub (pinned to 86d9c8fc54)