{"record":{"id":"cc2e743c3fd234a4","repo":"apache/iceberg","slug":"fail-to-deserialize-aggregated-statistics-change-t-cc2e74","errorCode":null,"errorMessage":"Fail to deserialize aggregated statistics,change to v1","messagePattern":"Fail to deserialize aggregated statistics,change to v1","errorType":"exception","errorClass":"RuntimeException","httpStatus":null,"severity":"error","filePath":"flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java","lineNumber":81,"sourceCode":"  static byte[] serializeCompletedStatistics(\n      CompletedStatistics completedStatistics,\n      TypeSerializer<CompletedStatistics> statisticsSerializer) {\n    try {\n      DataOutputSerializer out = new DataOutputSerializer(1024);\n      statisticsSerializer.serialize(completedStatistics, out);\n      return out.getCopyOfBuffer();\n    } catch (IOException e) {\n      throw new UncheckedIOException(\"Fail to serialize aggregated statistics\", e);\n    }\n  }\n\n  static CompletedStatistics deserializeCompletedStatistics(\n      byte[] bytes, CompletedStatisticsSerializer statisticsSerializer) {\n    try {\n      DataInputDeserializer input = new DataInputDeserializer(bytes);\n      CompletedStatistics completedStatistics = statisticsSerializer.deserialize(input);\n      if (!completedStatistics.isValid()) {\n        throw new RuntimeException(\"Fail to deserialize aggregated statistics,change to v1\");\n      }\n\n      return completedStatistics;\n    } catch (Exception e) {\n      try {\n        // If we restore from a lower version, the new version of SortKeySerializer cannot correctly\n        // parse the checkpointData, so we need to first switch the version to v1. Once the state\n        // data is successfully parsed, we need to switch the serialization version to the latest\n        // version to parse the subsequent data passed from the TM.\n        statisticsSerializer.changeSortKeySerializerVersion(1);\n        DataInputDeserializer input = new DataInputDeserializer(bytes);\n        CompletedStatistics deserialize = statisticsSerializer.deserialize(input);\n        statisticsSerializer.changeSortKeySerializerVersionLatest();\n        return deserialize;\n      } catch (IOException ioException) {\n        throw new UncheckedIOException(\"Fail to deserialize aggregated statistics\", ioException);\n      }\n    }","sourceCodeStart":63,"sourceCodeEnd":99,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java#L63-L99","documentation":"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.","triggerScenarios":"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.","commonSituations":"Restoring a Flink job from a checkpoint/savepoint created by an older Iceberg release (v1 sketch format) while running new code; mixed-version clusters.","solutions":["Ensure the job is restored with the same Iceberg version that created the checkpoint/savepoint, or rely on the built-in v1 fallback path.","Verify CompletedStatisticsSerializer's sort key serializer version matches the serialized data (latest vs v1).","Discard stale shuffle state and restart to regenerate statistics if migration is not needed.","Inspect the wrapped cause to determine whether data is genuinely corrupt versus version-mismatched."],"exampleFix":"// before\nCompletedStatistics stats = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer);\n// after (force v1 handling when restoring old state)\nserializer.changeSortKeySerializerVersion(1);\nCompletedStatistics stats = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer);\nserializer.changeSortKeySerializerVersionLatest();","handlingStrategy":"fallback","validationCode":"// detect old-format state before restore and pre-set the serializer version\nboolean legacyState = checkpointMeta.containsKey(\"iceberg.shuffle.statistics.v1\");\nif (legacyState) serializer.changeSortKeySerializerVersion(1);","typeGuard":null,"tryCatchPattern":"try {\n  CompletedStatistics s = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer);\n} catch (RuntimeException e) {\n  LOG.warn(\"invalid completed statistics after latest+fallback attempts; discarding\", e);\n  // proceed with empty statistics rather than failing the job\n}","preventionTips":["Check the release notes for shuffle statistics format changes before upgrading.","Test job restore from old savepoints in staging before production upgrade.","Avoid mixing cluster node versions while shuffle statistics are in flight."],"tags":["flink","deserialization","version-migration","shuffle"],"backgroundTag":"incompatible-source-type","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}