Repository navigation
Fix GroupValues retained memory accounting - #25188
Conversation
- GroupValues::size() charges owner + retained buffers. - Column: added group-index, emit, vectorized, and vec backing. - Added capacity/reuse tests. - Updated tight spill test pools.
…) calls - Updated `row.rs` to subtract inline `RowConverter`/`Rows` descriptors when calculating nested `.size()` calls. - The outer `GroupValuesRows` descriptor is now charged only once, preventing duplicate size accounting. - This resolves overcounting of size for nested rows, improving the accuracy of size calculations. - No functional changes to the row data itself; only the size accounting logic has been refined. - Enhances performance and reliability for operations that rely on precise size metrics.
- Added a dedicated integration test in `datafusion/physical-plan/src/aggregates/group_values/multi_group_by/mod.rs` - The test validates scratch‑capacity and reuse behavior across: - Five distinct vectorized buffers - The `emit_scratch` logic path - Verifies that buffer growth (delta) correctly allocates additional capacity - Confirms that clearing the buffers retains the allocated charge, ensuring proper reuse without unnecessary reallocations
…pe row‑backed inline nested descriptors, add primitive row‑backed multi‑column tests, and enable 1,024 spill limit - All missing GroupColumn owner descriptors charged. - Row‑backed inline nested descriptors deduped. - Added primitive, row‑backed, multi‑column tests. - Existing 1,024 spill limit now spills; ordered test passes.
- Added `size_retains_reusable_buffers_after_emit` test to verify reusable buffers are retained after emit. - Checks `rows_buffer` and hash scratch retention after `EmitTo::First`. - Re‑interns values and rechecks accounting to ensure correct behavior after re‑emission.
- Moved comment to correct emit test.
- Increase SLT peak from **9.2 KB** to **9.4 KB**. - Retain spill test plan while adding assertions for **spill count** and **spill bytes**. - Update pool sizes: **non‑distinct 1,000,000** entries and **DISTINCT 4,256,000** entries. - Revised overall plan to reflect the metric adjustments and new test assertions.
…TINCT pool - Modified the partial‑aggregation logic to skip only when memory limits are exceeded, rather than under broader conditions. - Reduced the DISTINCT pool size from `4_256_000` to `1_000_000` to lower memory consumption and improve performance in constrained environments.
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #25188 +/- ##
==========================================
+ Coverage 82.66% 82.72% +0.06%
==========================================
Files 1147 1147
Lines 446357 449087 +2730
Branches 446357 449087 +2730
==========================================
+ Hits 368971 371512 +2541
+ Misses 54997 54924 -73
- Partials 22389 22651 +262 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
…llability spill test - Updated `nested_nullability.rs` test budget from `1_000_000` to `4_256_000`. - Identified cause: DISTINCT struct state experiences a transient peak that exhausts the 1 MiB fair pool during spilling. - Adjusted the budget to accommodate the peak memory usage, ensuring the test passes reliably.
- Pins `target_partitions=1` for the DISTINCT spill scenario, ensuring a deterministic topology. - Retains the 1 MiB pool and associated spill assertions to validate correctness. - Removes scheduler‑dependent Partial/Final aggregate competition that was causing out‑of‑memory (OOM) failures.
…iation during spill recovery - After materializing spill state, rebuild accumulators from empty equivalents. - Releases retained accumulator capacity before reservation reconciliation. - Covers hash + ordered single/final spill paths.
- Updated `nested_nullability.rs` to wrap `FairSpillPool` in `TrackConsumersPool`. - Next OOM includes: - consumer names - spillability - live reservations - peak bytes - Fair allocation behavior remains unchanged.
…ll regression and adjust spill behavior - Disables `single_distinct_aggregation_to_group_by` only in the DISTINCT spill regression. - Asserts that the physical plan contains exactly one `AggregateExec`. - Retains the 1‑partition + 1 MiB fair spill pool configuration. - Leaves rewrite coverage elsewhere unchanged.
…st and rebuild empty partial table - Partial OOM drain now emits states via `EmitTo::First(batch_size)` instead of full materialization. - Eliminates full state‑batch materialization before slicing, reducing memory overhead. - Rebuilds an empty partial table after drain and resumes input processing to maintain continuity. - Addresses a 500 B grouping‑set regression and adds an `early_emit_count` assertion to catch premature emissions. - Updates partial hash tests to verify bounded incremental output behavior.
…r messages to match FinalHashAggregateStream
deb8299 to
dde38d4
Compare
…egate memory spill test
…dd regression test - Shrinks the group‑index and emits the scratch vector. - Releases and rebuilds the vector’s scratch capacity. - Resets the Boolean bitmap capacity safely. - Adds a zero‑reset regression test to verify the behavior.
…he ownership - `cached_values` changed from `Weak<dyn Array>` to `Option<Weak<dyn Array>>`. - Cache retains its identity hint; it no longer owns or pins the input allocation. - Added regression test to confirm a live identity cache while the input dictionary is released after drop. - Reused `clear()` in the shrink path, resolving warnings‑denied dead‑code lint.
…_by): shrink_to remaining_row_indices instead of clear Optimizes memory usage by calling `shrink_to(num_rows)` on `self.remaining_row_indices` rather than `clear()`. This avoids reallocating a new vector each time the buffer is reused, reducing overhead and improving performance in vectorized aggregations.
396eb7d to
cf77086
Compare
- Implemented descriptor-aware legacy spill budget logic. - Re-exported test-required GroupValuesPrimitive for broader test coverage. - Preserved existing spill, replay, and cleanup assertions.
- Removed broad post-state(EmitTo::All) accumulator rebuilds. - Kept legitimate fresh accumulators for partial-skip table.
- 2 expectations restored → PartialHashAggregateStream[0].
…BufferBuilder - Replaced finish/append/truncate with BooleanBufferBuilder::new(num_rows). - Removed obsolete comment.
- Added is_cached(&self, values) method. - Used std::ptr::addr_eq on Weak::as_ptr / Arc::as_ptr. - Replaced cache checks in sync_value_cache + tests. - Kept upgrade() only for weak-liveness test.
…port, update budget, clarify headroom - Removed test-only GroupValuesPrimitive re-export. - Replaced private-size-derived budget with MEMORY_LIMIT: 8_384. - Added concise aggregate-state headroom rationale.
…ry_spill.slt - Case G memory limit: 2M → 1M. - Existing result, final-spill, partial early-emission assertions retained.
|
8ec5f3e - #25734, #25735, #25736 - Follow-up idea, not for this PR: the maps in GroupValuesPrimitive, GroupValuesRows and GroupValuesColumn are still counted as capacity() * size_of::(). That is about 7/8 of the buckets and excludes control bytes. ArrowBytesMap::size() already uses HashTable::allocation_size(). |
…act `HashTable::allocation_size()` in `GroupValuesRows` (apache#25750) ## Which issue does this PR close? Closes apache#25735 Related to apache#25188 ## Rationale for this change **Background** `GroupValuesRows` is the general-purpose fallback implementation used for grouping complex nested types (like `Struct`, `List`, `Map`) and multi-column schemas. Its `size()` method is highly critical because it tells DataFusion's memory pool exactly how much memory is being consumed, which dictates when the system should safely spill to disk to prevent Out of Memory (OOM) errors. **The Problem** Previously, the memory size of the hash table was poorly estimated. It naively calculated the size as `capacity * size_of::<(u64, usize)>()`. This completely ignored the actual internal memory overhead of the hash table (like control bytes and entry slots), meaning DataFusion was under-reporting its memory usage. ## What changes are included in this PR? This PR completely removes the inaccurate manual `self.map_size` calculation. Instead, it directly calls `self.map.allocation_size()` inside the `size()` function. This ensures that DataFusion's memory pool accurately accounts for every single byte of overhead used by the underlying hash table. ## Are these changes tested? Yes! I have added two comprehensive tests to verify the exact memory accounting: - `test_exact_hash_table_allocation_accounting_multi_column`: Validates precise memory tracking during inserts, partial emits, full emits, and table shrink operations for multi-column groupings. - `test_exact_hash_table_allocation_accounting_nested_keys`: Validates the exact tracking for nested types (Lists) to ensure the row converter and map sizes are perfectly aligned. ## Are there any user-facing changes? No API changes. This is strictly an internal reliability fix to improve the accuracy of query execution memory tracking and prevent potential OOM crashes. Co-authored-by: kosiew <kosiew@gmail.com>
Upstream apache#25652 added Arc::ptr_eq checks against cached_values, which this branch changed to Option<Weak<dyn Array>>. Use the existing is_cached helper instead.
- **hash_stream.rs**: holder baseline now excludes the partition‑0 initial reservation. - **mod.rs**: set an exact initial‑table pool limit and removed the stale spill assertion.
|
🚀 |
- Update memory pool usage in hash stream to support unbounded limits by introducing UnboundedMemoryPool - Modify shared pool tests to use unbounded pools when expected to exceed reference budget - Ensure consistent memory pool selection based on input limit in test cases
Which issue does this PR close?
sizefunctions #23393Rationale for this change
Aggregate memory estimates can omit storage that group values keep after an operation or an emit. Under a memory limit, that can cause the aggregate’s reservation to understate what it retains and affect when it emits or spills.
What changes are included in this PR?
GroupValues::size(), counting reusable buffers by capacity and avoiding duplicate descriptor charges.Are these changes tested?
Yes. The patch adds capacity and emit/reuse tests for row and column group values, a dictionary-cache lifetime test, and assertions that the memory-limited nested-nullability cases actually spill. It also tightens an early-emission assertion, adjusts spill test limits, and updates an
EXPLAIN ANALYZEexpected peak-memory value.Are there any user-facing changes?
No public API change is shown. Under memory pressure, queries may emit partial aggregate results or spill at different points because retained memory is accounted for more fully.
LLM-generated code disclosure
This PR includes LLM-generated code and comments. All LLM-generated content has been manually reviewed.