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]
