{"record":{"id":"0cc9710bb81745f6","repo":"apache/iceberg","slug":"fail-to-serialize-data-statistics-0cc971","errorCode":null,"errorMessage":"Fail to serialize data statistics","messagePattern":"Fail to serialize data 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":49,"sourceCode":"\n  static DataStatistics createTaskStatistics(\n      StatisticsType type, int operatorParallelism, int numPartitions) {\n    if (type == StatisticsType.Map) {\n      return new MapDataStatistics();\n    } else {\n      return new SketchDataStatistics(\n          SketchUtil.determineOperatorReservoirSize(operatorParallelism, numPartitions));\n    }\n  }\n\n  static byte[] serializeDataStatistics(\n      DataStatistics dataStatistics, TypeSerializer<DataStatistics> statisticsSerializer) {\n    DataOutputSerializer out = new DataOutputSerializer(64);\n    try {\n      statisticsSerializer.serialize(dataStatistics, out);\n      return out.getCopyOfBuffer();\n    } catch (IOException e) {\n      throw new UncheckedIOException(\"Fail to serialize data statistics\", e);\n    }\n  }\n\n  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);","sourceCodeStart":31,"sourceCodeEnd":67,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java#L31-L67","documentation":"Thrown by StatisticsUtil.serializeDataStatistics when the Flink TypeSerializer cannot serialize a DataStatistics object to bytes, wrapping the IOException in an UncheckedIOException. Data statistics are serialized to be carried through Flink accumulators/state during the shuffle-based sort. Any IO failure in the serializer's internal write path triggers this.","triggerScenarios":"Calling serializeDataStatistics with a DataStatistics whose underlying serializer (e.g. the map-based LocalStatisticsSerializer) hits an IOException, typically from DataOutputSerializer write failures or a null/incompatible statistics object.","commonSituations":"Passing a DataStatistics of a type not matching the supplied TypeSerializer after a code change, or serializer bugs when emitting shuffle statistics from a downstream subtask to the coordinator.","solutions":["Verify the TypeSerializer passed in matches the concrete DataStatistics implementation being serialized.","Upgrade to a fixed Iceberg version if the serializer itself is buggy (check issue tracker).","Add a guard that statistics are non-null and of the expected class before calling serializeDataStatistics."],"exampleFix":"// before\nbyte[] bytes = StatisticsUtil.serializeDataStatistics(stats, new LocalStatisticsSerializer());\n// after\nPreconditions.checkArgument(stats instanceof LocalStatistics,\n    \"Expected LocalStatistics, got %s\", stats.getClass());\nbyte[] bytes = StatisticsUtil.serializeDataStatistics(stats, new LocalStatisticsSerializer());","handlingStrategy":"type-guard","validationCode":"if (dataStatistics == null || statisticsSerializer == null) { throw new IllegalArgumentException(\"statistics and serializer required\"); }","typeGuard":"boolean isSerializable(DataStatistics s, TypeSerializer<DataStatistics> ser) { return s != null && ser != null; }","tryCatchPattern":"try {\n  byte[] bytes = StatisticsUtil.serializeDataStatistics(stats, serializer);\n} catch (UncheckedIOException e) {\n  LOG.error(\"Statistics serialization failed\", e);\n  throw e; // statistics loss is not recoverable locally\n}","preventionTips":["Always pair a DataStatistics implementation with its matching serializer type.","Keep serializer and statistics classes in sync across Iceberg upgrades.","Unit-test round-trip serialize/deserialize after any serializer change."],"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"}