Conversation
…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.
sunchao
left a comment
There was a problem hiding this comment.
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 smallheap_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.
Which issue does this PR close?
Closes #6252.
Rationale for this change
SumIntGroupsAccumulatorLegacy,SumIntGroupsAccumulatorAnsiandSumIntGroupsAccumulatorTryinnative/spark-expr/src/agg_funcs/sum_int.rsreturnedstd::mem::size_of_val(self)fromGroupsAccumulator::size(). That is the size of the struct, which holds only theVecheaders, so the per-group state was left out: 16 bytes per group forsums: Vec<Option<i64>>, plus 1 byte per group forhas_all_nullsin theTryvariant. DataFusion sizes a grouped hash aggregate's memory reservation from its accumulators'size()(AggregateHashTable::memory_sizein datafusion-physical-plan 55.1.0), so every groupedSUMover anInt8,Int16,Int32orInt64column reserved far less than it held and spilled later than it should. With 1M groups, each accumulator reported 24 bytes (48 forTry) instead of at least 16 MB.I audited the other
AccumulatorandGroupsAccumulatorimpls in the native crates for the same mistake and found one more: theAccumulatorimpl forSparkBloomFilterinnative/spark-expr/src/bloom_filter/bloom_filter_agg.rsalso returnedsize_of_val(self), leaving out the bit array the filter allocates up front (num_bits / 8bytes, up to 8 MiB under Spark's defaultspark.sql.optimizer.runtime.bloomFilter.maxNumBits). Its effect today is smaller. Spark plansBloomFilterAggregatefor runtime filters as an ungrouped aggregate, and DataFusion'sAggregateStreamonly reserves the growth insize()acrossupdate_batchandmerge_batch, which is zero for a filter that never resizes. The change makessize()follow theAccumulatorcontract, which matters wherever the full size is read, for exampleGroupsAccumulatorAdapter, which counts each wrapped accumulator'ssize()when it creates it.The remaining
size_of_val(self)returns are on scalar accumulators whose fields are all inline (the threeSumIntegerAccumulator*types,AvgAccumulator,AvgDecimalAccumulator,SumDecimalAccumulator,VarianceAccumulatorandCovarianceAccumulator), 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()andSumIntGroupsAccumulatorAnsi::size()return the capacity ofsumstimessize_of::<Option<i64>>(), andSumIntGroupsAccumulatorTry::size()also adds the capacity ofhas_all_nulls. This mirrorsSumDecimalGroupsAccumulator::size().Accumulator::size()forSparkBloomFilteradds the bytes allocated for its bit array, read through newheap_size()methods onSparkBloomFilterandSparkBitArray.How are these changes tested?
New Rust unit tests.
test_legacy_size_counts_group_state,test_ansi_size_counts_group_stateandtest_try_size_counts_group_stateinsum_int.rsrunupdate_batchover 1M distinct groups and assert thatsize()is at least 1M times the per-group state (16 bytes, or 17 bytes forTry).size_counts_bit_arrayinbloom_filter_agg.rsasserts 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
Tryaccumulator 48 bytes. With the change,cargo test -p datafusion-comet-spark-expr --lib -- sum_int bloom_filterpasses (37 tests).cargo fmt --allandcargo clippy --all-targets --workspace -- -D warningsare clean. No JVM code changed.