{"record":{"id":"0ad3c4f403b95d4d","repo":"apache/iceberg","slug":"fail-to-serialize-data-statistics-0ad3c4","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.3/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.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java#L31-L67","documentation":"StatisticsUtil.serializeDataStatistics wraps any IOException thrown while serializing a DataStatistics object (used in Flink key shuffle downsketching) into an UncheckedIOException. This is an internal failure of the Flink TypeSerializer for the DataStatistics type; the library throws it because serialization failures cannot propagate as checked IOExceptions through non-declaring call sites.","triggerScenarios":"Calling StatisticsUtil.serializeDataStatistics(dataStatistics, serializer) when the provided TypeSerializer<DataStatistics> throws IOException during serialize(), e.g. because the serializer's internal buffer cannot grow or the object graph is inconsistent with the serializer version.","commonSituations":"Job upgrade/downgrade where the serialized statistics type no longer matches the serializer handed to the utility; a custom or misregistered TypeSerializer that fails mid-serialization; JVM heap/buffer issues during state transfer in the shuffle.","solutions":["Verify the TypeSerializer<DataStatistics> passed in matches the DataStatistics implementation currently in use (SketchDataStatistics vs MapDataStatistics) and its version.","Check for Flink/Iceberg version mismatch between job manager and task manager classpaths; align all nodes to the same iceberg-flink-runtime version.","Reproduce the wrapped cause via e.getCause() to see the underlying IOException and fix the serializer accordingly.","If a job restore triggered it, restart the job without state or migrate state with a matching serializer version."],"exampleFix":"// before\nbyte[] bytes = StatisticsUtil.serializeDataStatistics(stats, wrongSerializer);\n// after\nTypeSerializer<DataStatistics> serializer =\n    new DataStatisticsSerializer(); // matches DataStatistics impl and version\nbyte[] bytes = StatisticsUtil.serializeDataStatistics(stats, serializer);","handlingStrategy":"try-catch","validationCode":"if (dataStatistics == null) throw new IllegalArgumentException(\"dataStatistics must not be null\");\n// ensure serializer type matches the DataStatistics impl before calling\nPreconditions.checkArgument(serializer instanceof DataStatisticsSerializer);","typeGuard":"boolean isCompatible(TypeSerializer<DataStatistics> s, DataStatistics d) {\n  return s != null && d != null;\n}","tryCatchPattern":"try {\n  byte[] bytes = StatisticsUtil.serializeDataStatistics(stats, serializer);\n} catch (UncheckedIOException e) {\n  IOException cause = e.getCause();\n  LOG.error(\"statistics serialization failed\", cause);\n  throw e; // or fall back to empty statistics\n}","preventionTips":["Pin the same iceberg-flink-runtime version on all task/job managers.","Only use serializers obtained from the matching DataStatistics type.","Log e.getCause() to diagnose serializer version drift early."],"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"}