apache/iceberg · error · RuntimeException
Fail to deserialize aggregated statistics,change to v1
Error message
Fail to deserialize aggregated statistics,change to v1
What it means
StatisticsUtil.deserializeCompletedStatistics throws this RuntimeException when deserialized CompletedStatistics fails its isValid() check, meaning the bytes did not yield a valid v2 (latest-version) aggregate statistics payload. Unlike the other methods, it first attempts a v1 fallback by switching the SortKeySerializer to version 1.
Solutions
- Ensure the job is restored with the same Iceberg version that created the checkpoint/savepoint, or rely on the built-in v1 fallback path.
- Verify CompletedStatisticsSerializer's sort key serializer version matches the serialized data (latest vs v1).
- Discard stale shuffle state and restart to regenerate statistics if migration is not needed.
- Inspect the wrapped cause to determine whether data is genuinely corrupt versus version-mismatched.
Example fix
// before CompletedStatistics stats = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer); // after (force v1 handling when restoring old state) serializer.changeSortKeySerializerVersion(1); CompletedStatistics stats = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer); serializer.changeSortKeySerializerVersionLatest();
Defensive patterns
Strategy: fallback
Validate before calling
// detect old-format state before restore and pre-set the serializer version
boolean legacyState = checkpointMeta.containsKey("iceberg.shuffle.statistics.v1");
if (legacyState) serializer.changeSortKeySerializerVersion(1); Try / catch
try {
CompletedStatistics s = StatisticsUtil.deserializeCompletedStatistics(bytes, serializer);
} catch (RuntimeException e) {
LOG.warn("invalid completed statistics after latest+fallback attempts; discarding", e);
// proceed with empty statistics rather than failing the job
} Prevention
- Check the release notes for shuffle statistics format changes before upgrading.
- Test job restore from old savepoints in staging before production upgrade.
- Avoid mixing cluster node versions while shuffle statistics are in flight.
When it happens
Trigger: Calling StatisticsUtil.deserializeCompletedStatistics(bytes, serializer) where the serializer's deserialize() returns an invalid CompletedStatistics (isValid() == false) — typically bytes written with a different statistics version than the current serializer configuration.
Common situations: Restoring a Flink job from a checkpoint/savepoint created by an older Iceberg release (v1 sketch format) while running new code; mixed-version clusters.
Understand the failure class
Background: "is not a compatible type" / "cannot merge" errors: when a value's type doesn't match what the library requires — this error's family across 65 libraries.
Related errors
- Fail to deserialize data statistics
- Fail to deserialize aggregated statistics
- Fail to deserialize aggregated statistics
- Fail to deserialize aggregated statistics,change to v1
- Fail to deserialize aggregated statistics,change to v1
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/cc2e743c3fd234a4.
Report an issue: GitHub.
Appendix: source
Thrown at flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/StatisticsUtil.java:81
static byte[] serializeCompletedStatistics(
CompletedStatistics completedStatistics,
TypeSerializer<CompletedStatistics> statisticsSerializer) {
try {
DataOutputSerializer out = new DataOutputSerializer(1024);
statisticsSerializer.serialize(completedStatistics, out);
return out.getCopyOfBuffer();
} catch (IOException e) {
throw new UncheckedIOException("Fail to serialize aggregated statistics", e);
}
}
static CompletedStatistics deserializeCompletedStatistics(
byte[] bytes, CompletedStatisticsSerializer statisticsSerializer) {
try {
DataInputDeserializer input = new DataInputDeserializer(bytes);
CompletedStatistics completedStatistics = statisticsSerializer.deserialize(input);
if (!completedStatistics.isValid()) {
throw new RuntimeException("Fail to deserialize aggregated statistics,change to v1");
}
return completedStatistics;
} catch (Exception e) {
try {
// If we restore from a lower version, the new version of SortKeySerializer cannot correctly
// parse the checkpointData, so we need to first switch the version to v1. Once the state
// data is successfully parsed, we need to switch the serialization version to the latest
// version to parse the subsequent data passed from the TM.
statisticsSerializer.changeSortKeySerializerVersion(1);
DataInputDeserializer input = new DataInputDeserializer(bytes);
CompletedStatistics deserialize = statisticsSerializer.deserialize(input);
statisticsSerializer.changeSortKeySerializerVersionLatest();
return deserialize;
} catch (IOException ioException) {
throw new UncheckedIOException("Fail to deserialize aggregated statistics", ioException);
}
}View on GitHub (pinned to 86d9c8fc54)