{"record":{"id":"e1572a028371d795","repo":"apache/iceberg","slug":"fail-to-serialize-aggregated-statistics-e1572a","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.3/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.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java#L53-L89","documentation":"StatisticsUtil.serializeCompletedStatistics wraps IOException thrown while serializing a CompletedStatistics object (aggregated shuffle statistics sent from task to coordinator) into an UncheckedIOException. It signals that the CompletedStatisticsSerializer failed to write the aggregate payload.","triggerScenarios":"Calling StatisticsUtil.serializeCompletedStatistics(completedStatistics, statisticsSerializer) when statisticsSerializer.serialize() throws IOException, e.g. serializer version mismatch with the contained statistics entries or buffer allocation failure.","commonSituations":"Coordinator/task communication after an Iceberg upgrade where the SortKeySerializer version no longer matches; custom serializers not registered on all workers.","solutions":["Confirm CompletedStatistics and its inner DataStatistics match the serializer version (Sketch vs Map, sort key version).","Align iceberg-flink-runtime version on all task and job manager nodes.","Get the root cause from e.getCause() and fix or upgrade the serializer.","Fall back to StatisticsType.Map statistics if sketch serialization keeps failing after version changes."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"Preconditions.checkArgument(completedStatistics != null, \"completedStatistics must not be null\");\nPreconditions.checkArgument(statisticsSerializer != null, \"statisticsSerializer must not be null\");","typeGuard":null,"tryCatchPattern":"try {\n  byte[] bytes = StatisticsUtil.serializeCompletedStatistics(stats, serializer);\n} catch (UncheckedIOException e) {\n  LOG.error(\"completed statistics serialization failed\", e.getCause());\n  throw e;\n}","preventionTips":["Keep CompletedStatisticsSerializer sort key version consistent with the statistics payload.","Upgrade all cluster nodes together; avoid rolling upgrades across Iceberg majors.","Inspect cause chains when serialization fails after a version bump."],"tags":["flink","serialization","unchecked-io","shuffle"],"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"}