{"record":{"id":"6c864b1765c02f3a","repo":"apache/iceberg","slug":"fail-to-deserialize-data-statistics-6c864b","errorCode":null,"errorMessage":"Fail to deserialize data statistics","messagePattern":"Fail to deserialize 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":59,"sourceCode":"\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);\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 {","sourceCodeStart":41,"sourceCodeEnd":77,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java#L41-L77","documentation":"Thrown by StatisticsUtil.deserializeDataStatistics when bytes cannot be decoded back into a DataStatistics via the given TypeSerializer; the IOException is wrapped as UncheckedIOException. This occurs when reading shuffle statistics persisted in Flink state or transmitted between operators. Mismatched serializer versions or corrupted bytes cause it.","triggerScenarios":"Calling deserializeDataStatistics on bytes produced by a different serializer version, truncated byte arrays, or deserializing statistics written by another Iceberg release.","commonSituations":"Job restore from checkpoints across Iceberg upgrades, mixing serializer versions between writer and reader subtasks during a rolling upgrade.","solutions":["Restore the job with the same Iceberg version that created the state, or migrate state deliberately.","Confirm the serializer passed matches the one used at serialization time (statistics type must match).","Discard stale state and restart the job so statistics are re-collected from scratch."],"exampleFix":"// before: bytes from old job version\nDataStatistics stats = StatisticsUtil.deserializeDataStatistics(bytes, serializer);\n// after: null/empty guard and version-validated path\nif (bytes == null || bytes.length == 0) return DataStatistics.empty();\nDataStatistics stats = StatisticsUtil.deserializeDataStatistics(bytes, serializer);","handlingStrategy":"validation","validationCode":"if (bytes == null || bytes.length == 0) { return null; } // skip empty payloads before deserializing","typeGuard":"boolean hasPayload(byte[] bytes) { return bytes != null && bytes.length > 0; }","tryCatchPattern":"try {\n  DataStatistics s = StatisticsUtil.deserializeDataStatistics(bytes, serializer);\n} catch (UncheckedIOException e) {\n  LOG.warn(\"Cannot decode statistics; treating as empty\", e);\n  DataStatistics s = DataStatistics.empty();\n}","preventionTips":["Restore Flink state with the Iceberg version that produced it.","Avoid rolling upgrades that mix serializer versions between subtasks.","Persist serializer version metadata alongside statistics bytes."],"tags":["serialization","flink","unchecked-io"],"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"}