{"record":{"id":"3c7f7342c31e7d25","repo":"apache/iceberg","slug":"fail-to-deserialize-aggregated-statistics-3c7f73","errorCode":null,"errorMessage":"Fail to deserialize aggregated statistics","messagePattern":"Fail to deserialize aggregated statistics","errorType":"exception","errorClass":"UncheckedIOException","httpStatus":null,"severity":"error","filePath":"flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java","lineNumber":97,"sourceCode":"      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    }\n  }\n\n  static byte[] serializeGlobalStatistics(\n      GlobalStatistics globalStatistics, TypeSerializer<GlobalStatistics> statisticsSerializer) {\n    try {\n      DataOutputSerializer out = new DataOutputSerializer(1024);\n      statisticsSerializer.serialize(globalStatistics, out);\n      return out.getCopyOfBuffer();\n    } catch (IOException e) {\n      throw new UncheckedIOException(\"Fail to serialize aggregated statistics\", e);\n    }\n  }\n\n  static GlobalStatistics deserializeGlobalStatistics(\n      byte[] bytes, TypeSerializer<GlobalStatistics> statisticsSerializer) {\n    try {","sourceCodeStart":79,"sourceCodeEnd":115,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java#L79-L115","documentation":"StatisticsUtil.deserializeCompletedStatistics throws this UncheckedIOException when even the v1 fallback deserialization (after changeSortKeySerializerVersion(1)) fails with an IOException, meaning the bytes are unreadable under any known statistics format version.","triggerScenarios":"Calling StatisticsUtil.deserializeCompletedStatistics where both the latest-version attempt (invalid result) and the v1 retry path throw IOException — bytes are corrupt or from an incompatible/unrecognized format.","commonSituations":"Truncated or corrupted checkpoint data; statistics bytes produced by a much older or newer format with no compatible reader; manual state surgery.","solutions":["Treat the state as unrecoverable: discard the checkpoint/savepoint and restart the job to regenerate shuffle statistics.","Verify checkpoint storage integrity (DFS corruption, incomplete upload) and check TaskManager logs for the original IOException cause.","Ensure all nodes run the identical iceberg-flink-runtime version; a mixed cluster can produce unreadable aggregates.","If the issue follows an upgrade, downgrade to the version that wrote the state, then migrate forward cleanly."],"exampleFix":null,"handlingStrategy":"fallback","validationCode":"// verify checkpoint storage integrity before restore\nFileSystem fs = FileSystem.get(checkpointUri, conf);\nif (!fs.exists(new Path(checkpointUri, \"_metadata\"))) throw new IllegalStateException(\"incomplete checkpoint\");","typeGuard":null,"tryCatchPattern":"try {\n  CompletedStatistics s = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer);\n} catch (UncheckedIOException e) {\n  LOG.error(\"statistics unreadable under any format version; regenerating\", e.getCause());\n  // fall back to empty statistics / restart without state\n}","preventionTips":["Enable checkpoint verification/monitoring to catch truncated DFS uploads.","Do not manually edit or trim operator state files.","Restart from a known-good savepoint when deserialization of state fails."],"tags":["flink","deserialization","unchecked-io","corruption"],"backgroundTag":"json-unmarshal-failed","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"}