apache/iceberg · error · RuntimeException

Fail to deserialize aggregated statistics,change to v1

Error message

Fail to deserialize aggregated statistics,change to v1

What it means

StatisticsUtil.deserializeCompletedStatistics throws this RuntimeException when deserialized CompletedStatistics fails its isValid() check, meaning the bytes did not yield a valid v2 (latest-version) aggregate statistics payload. Unlike the other methods, it first attempts a v1 fallback by switching the SortKeySerializer to version 1.

Solutions

  1. Ensure the job is restored with the same Iceberg version that created the checkpoint/savepoint, or rely on the built-in v1 fallback path.
  2. Verify CompletedStatisticsSerializer's sort key serializer version matches the serialized data (latest vs v1).
  3. Discard stale shuffle state and restart to regenerate statistics if migration is not needed.
  4. Inspect the wrapped cause to determine whether data is genuinely corrupt versus version-mismatched.

Example fix

// before
CompletedStatistics stats = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer);
// after (force v1 handling when restoring old state)
serializer.changeSortKeySerializerVersion(1);
CompletedStatistics stats = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer);
serializer.changeSortKeySerializerVersionLatest();
Defensive patterns

Strategy: fallback

Validate before calling

// detect old-format state before restore and pre-set the serializer version
boolean legacyState = checkpointMeta.containsKey("iceberg.shuffle.statistics.v1");
if (legacyState) serializer.changeSortKeySerializerVersion(1);

Try / catch

try {
  CompletedStatistics s = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer);
} catch (RuntimeException e) {
  LOG.warn("invalid completed statistics after latest+fallback attempts; discarding", e);
  // proceed with empty statistics rather than failing the job
}

Prevention

When it happens

Trigger: Calling StatisticsUtil.deserializeCompletedStatistics(bytes, serializer) where the serializer's deserialize() returns an invalid CompletedStatistics (isValid() == false) — typically bytes written with a different statistics version than the current serializer configuration.

Common situations: Restoring a Flink job from a checkpoint/savepoint created by an older Iceberg release (v1 sketch format) while running new code; mixed-version clusters.

Understand the failure class

Background: "is not a compatible type" / "cannot merge" errors: when a value's type doesn't match what the library requires — this error's family across 65 libraries.

Related errors


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

Appendix: source

Thrown at flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java:81

  static byte[] serializeCompletedStatistics(
      CompletedStatistics completedStatistics,
      TypeSerializer<CompletedStatistics> statisticsSerializer) {
    try {
      DataOutputSerializer out = new DataOutputSerializer(1024);
      statisticsSerializer.serialize(completedStatistics, out);
      return out.getCopyOfBuffer();
    } catch (IOException e) {
      throw new UncheckedIOException("Fail to serialize aggregated statistics", e);
    }
  }

  static CompletedStatistics deserializeCompletedStatistics(
      byte[] bytes, CompletedStatisticsSerializer statisticsSerializer) {
    try {
      DataInputDeserializer input = new DataInputDeserializer(bytes);
      CompletedStatistics completedStatistics = statisticsSerializer.deserialize(input);
      if (!completedStatistics.isValid()) {
        throw new RuntimeException("Fail to deserialize aggregated statistics,change to v1");
      }

      return completedStatistics;
    } catch (Exception e) {
      try {
        // If we restore from a lower version, the new version of SortKeySerializer cannot correctly
        // parse the checkpointData, so we need to first switch the version to v1. Once the state
        // data is successfully parsed, we need to switch the serialization version to the latest
        // version to parse the subsequent data passed from the TM.
        statisticsSerializer.changeSortKeySerializerVersion(1);
        DataInputDeserializer input = new DataInputDeserializer(bytes);
        CompletedStatistics deserialize = statisticsSerializer.deserialize(input);
        statisticsSerializer.changeSortKeySerializerVersionLatest();
        return deserialize;
      } catch (IOException ioException) {
        throw new UncheckedIOException("Fail to deserialize aggregated statistics", ioException);
      }
    }

View on GitHub (pinned to 86d9c8fc54)