{"record":{"id":"4a62e7e2d9384347","repo":"apache/iceberg","slug":"fail-to-serialize-aggregated-statistics-4a62e7","errorCode":null,"errorMessage":"Fail to serialize aggregated statistics","messagePattern":"Fail to serialize aggregated statistics","errorType":"exception","errorClass":"UncheckedIOException","httpStatus":null,"severity":"error","filePath":"flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java","lineNumber":71,"sourceCode":"  static DataStatistics deserializeDataStatistics(\n      byte[] bytes, TypeSerializer<DataStatistics> statisticsSerializer) {\n    DataInputDeserializer input = new DataInputDeserializer(bytes, 0, bytes.length);\n    try {\n      return statisticsSerializer.deserialize(input);\n    } catch (IOException e) {\n      throw new UncheckedIOException(\"Fail to deserialize data statistics\", e);\n    }\n  }\n\n  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","sourceCodeStart":53,"sourceCodeEnd":89,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java#L53-L89","documentation":"Thrown by StatisticsUtil.serializeCompletedStatistics when serializing a CompletedStatistics (aggregated statistics from the coordinator) fails with an IOException. Completed statistics are persisted with a 1024-byte initial DataOutputSerializer and handed to state backends / committed metadata. The library wraps the IO failure as an UncheckedIOException.","triggerScenarios":"Calling serializeCompletedStatistics with the wrong CompletedStatisticsSerializer for the concrete statistics class, or an internal serializer IOException while writing aggregated sketch/map statistics.","commonSituations":"Coordinator persisting aggregated statistics at checkpoint time; mismatched serializer after upgrading Iceberg versions where the statistics layout changed.","solutions":["Ensure the CompletedStatisticsSerializer variant matches the statistics version produced by the aggregator.","Align Iceberg versions between checkpoint writing and restoring jobs.","Verify the completedStatistics object is non-null and valid (isValid()) before serializing."],"exampleFix":"// before\nbyte[] bytes = StatisticsUtil.serializeCompletedStatistics(completed, statisticsSerializer);\n// after\nif (completed == null || !completed.isValid()) {\n  throw new IllegalStateException(\"Completed statistics missing or invalid\");\n}\nbyte[] bytes = StatisticsUtil.serializeCompletedStatistics(completed, statisticsSerializer);","handlingStrategy":"validation","validationCode":"if (completedStatistics == null || !completedStatistics.isValid()) { throw new IllegalStateException(\"Completed statistics invalid before serialize\"); }","typeGuard":"boolean canSerialize(CompletedStatistics cs) { return cs != null && cs.isValid(); }","tryCatchPattern":"try {\n  byte[] bytes = StatisticsUtil.serializeCompletedStatistics(completed, serializer);\n} catch (UncheckedIOException e) {\n  LOG.error(\"Failed to persist aggregated statistics\", e);\n  throw e;\n}","preventionTips":["Validate CompletedStatistics with isValid() before serializing.","Match the CompletedStatisticsSerializer version to the aggregator's output version.","Test checkpoint/restore round trips in CI across Iceberg upgrades."],"tags":["serialization","flink","unchecked-io"],"backgroundTag":"json-serialization-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"}