Skip to content

fix: count grouped integer SUM state in the aggregate's memory reservation - #6370

Open
andygrove wants to merge 1 commit into
apache:mainfrom
andygrove:sum-int-accumulator-size
Open

andygrove wants to merge 1 commit into
apache:mainfrom
andygrove:sum-int-accumulator-size

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6252.

Rationale for this change

SumIntGroupsAccumulatorLegacy, SumIntGroupsAccumulatorAnsi and SumIntGroupsAccumulatorTry in native/spark-expr/src/agg_funcs/sum_int.rs returned std::mem::size_of_val(self) from GroupsAccumulator::size(). That is the size of the struct, which holds only the Vec headers, so the per-group state was left out: 16 bytes per group for sums: Vec<Option<i64>>, plus 1 byte per group for has_all_nulls in the Try variant. DataFusion sizes a grouped hash aggregate's memory reservation from its accumulators' size() (AggregateHashTable::memory_size in datafusion-physical-plan 55.1.0), so every grouped SUM over an Int8, Int16, Int32 or Int64 column reserved far less than it held and spilled later than it should. With 1M groups, each accumulator reported 24 bytes (48 for Try) instead of at least 16 MB.

I audited the other Accumulator and GroupsAccumulator impls in the native crates for the same mistake and found one more: the Accumulator impl for SparkBloomFilter in native/spark-expr/src/bloom_filter/bloom_filter_agg.rs also returned size_of_val(self), leaving out the bit array the filter allocates up front (num_bits / 8 bytes, up to 8 MiB under Spark's default spark.sql.optimizer.runtime.bloomFilter.maxNumBits). Its effect today is smaller. Spark plans BloomFilterAggregate for runtime filters as an ungrouped aggregate, and DataFusion's AggregateStream only reserves the growth in size() across update_batch and merge_batch, which is zero for a filter that never resizes. The change makes size() follow the Accumulator contract, which matters wherever the full size is read, for example GroupsAccumulatorAdapter, which counts each wrapped accumulator's size() when it creates it.

The remaining size_of_val(self) returns are on scalar accumulators whose fields are all inline (the three SumIntegerAccumulator* types, AvgAccumulator, AvgDecimalAccumulator, SumDecimalAccumulator, VarianceAccumulator and CovarianceAccumulator), so they are left as they are. The other grouped accumulators already count their heap state.

What changes are included in this PR?

  • SumIntGroupsAccumulatorLegacy::size() and SumIntGroupsAccumulatorAnsi::size() return the capacity of sums times size_of::<Option<i64>>(), and SumIntGroupsAccumulatorTry::size() also adds the capacity of has_all_nulls. This mirrors SumDecimalGroupsAccumulator::size().
  • Accumulator::size() for SparkBloomFilter adds the bytes allocated for its bit array, read through new heap_size() methods on SparkBloomFilter and SparkBitArray.
  • Unit tests for the three integer sum variants and the bloom filter.

How are these changes tested?

New Rust unit tests. test_legacy_size_counts_group_state, test_ansi_size_counts_group_state and test_try_size_counts_group_state in sum_int.rs run update_batch over 1M distinct groups and assert that size() is at least 1M times the per-group state (16 bytes, or 17 bytes for Try). size_counts_bit_array in bloom_filter_agg.rs asserts that a filter with 1M bits reports at least 128 KiB.

I ran the new tests before changing the accumulators and all four failed: the Legacy and Ansi accumulators reported 24 bytes for 1M groups and the Try accumulator 48 bytes. With the change, cargo test -p datafusion-comet-spark-expr --lib -- sum_int bloom_filter passes (37 tests). cargo fmt --all and cargo clippy --all-targets --workspace -- -D warnings are clean. No JVM code changed.

…ation

The three SumInt*GroupsAccumulator impls returned size_of_val(self) from
GroupsAccumulator::size(), which covers the Vec headers but not the
per-group sums (16 bytes per group, plus has_all_nulls for Try).
DataFusion sizes a hash aggregate's reservation from these sizes, so
grouped integer sums reserved far less than they held. Count the
capacity of the state vectors, as SumDecimalGroupsAccumulator does.

The Accumulator impl for SparkBloomFilter had the same mistake and left
out its bit array, so include that too.
@github-actions github-actions Bot added bug Something isn't working area:aggregation Hash aggregates, aggregate expressions area:expressions Expression evaluation labels Sep 29, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Grouped integer SUM reported only vector headers, omitting per-group allocations. Bloom filter accumulators omitted their allocated bit buffers.
  • Design approach: Calculate owned heap storage from vector capacities and element sizes, feeding DataFusion’s existing reservation machinery.
  • Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Checked allocation ownership, state emission, reservation updates, and Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. The diff preserves arithmetic, null handling, overflow behavior, and Bloom filter serialization.
  • Key design decisions: Using capacity() accounts for retained allocations. The calculations are constant-time and introduce no allocations or per-group scans. The small heap_size() helpers keep buffer accounting with its owner.
  • Implementation sketch: Update the three grouped SUM size() implementations, add Bloom filter heap accounting through two helpers, and add four regression tests.
  • Behavioral changes worth calling out: Grouped aggregates now reserve their previously omitted storage and can emit or spill earlier under memory pressure. The Bloom filter change corrects its reported size, but does not address DataFusion’s existing growth-only accounting for ungrouped aggregates.
  • Suggested improvements: None meeting the requested P1/P2 evidence bar.

Reviewed the entire four-file diff from 8805e3f066440752ac98d5b48e197b43e4fb9ebd to e4c80f7b217cdb1bf7ca12bd3c4ad29718c42aa2. Confirmed non-draft status. The snapshot and live discussion checks contained no reviews or comments, so no existing blockers were identified. Routed skills: review-comet-pr, review-comet-expression-pr, and review-comet-memory-pr.

Exact-head CI: Required Checks, native build, Rust tests, Spark 4.1 Comet suites, and TPC-H/TPC-DS checks passed. The Rust log confirms all four new regression tests passed within 1,814 passing tests. CI’s generated merge commit has the same tree as the reviewed head. Spark SQL, Iceberg, and macOS suites were skipped.

Validation: A local harness compiled the changed source files against cached DataFusion 55.1.0 and Arrow 59.3.0 dependencies and passed 30 tests. All four new regression tests failed with the base implementations restored. The normal local Cargo run could not resolve locked dependency hdfs-sys 0.3.1. The harness was not a clean workspace build, and no local JVM, Spark SQL, or forced-spill integration suite was run.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:aggregation Hash aggregates, aggregate expressions area:expressions Expression evaluation bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Grouped integer SUM doesn't count its per-group state in the aggregate's memory reservation

2 participants