andygrove opened a new pull request, #6370:
URL: https://github.com/apache/datafusion-comet/pull/6370

   ## 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.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to