Rachelint opened a new issue, #26124:
URL: https://github.com/apache/datafusion/issues/26124

   ### Is your feature request related to a problem or challenge?
   
   I'd like to track proactive flushing for Partial hash aggregation as a 
separate performance improvement.
   
   Keeping a larger Partial aggregation table does not always pay off: the 
additional reduction in intermediate rows may not offset the cost of 
maintaining a larger working set. Recent experiments suggest that flushing 
earlier can improve query performance while reducing the amount of state 
retained by Partial aggregation.
   
   It would be useful to explore this independently of Final aggregation 
bucketing, so we can understand the benefits and trade-offs of each 
optimization.
   
   ### Describe the solution you'd like
   
   Allow Partial aggregation to emit its intermediate states proactively when a 
suitable threshold is reached, then continue aggregating subsequent input 
batches.
   
   Some questions worth investigating:
   
   - **Flush policy:** Compare distinct-group-count thresholds, thresholds tied 
to the target batch size, and byte thresholds. Include group keys, hash-table 
storage, and accumulator state when evaluating memory usage. A group-count 
limit alone does not bound variable-length payloads or growing aggregate states.
   - **Interaction with skipping Partial aggregation:** Preserve meaningful 
reduction estimates across flushes and account for repeated groups appearing in 
different flush windows. Frequent flushing and skipping aggregation have 
different costs.
   - **Output and downstream costs:** Evaluate emitting the complete flushed 
batch versus slicing it, including the effects on repartitioning, coalescing, 
and Final aggregation.
   - **Allocation overhead:** Investigate capacity reuse and reservation as 
follow-ups, with controls that separate their effects from the flushing policy 
itself.
   
   The goal would be to find a policy that improves end-to-end performance 
without introducing substantial regressions for workloads where Partial 
aggregation already reduces the input effectively. Validation should cover 
integer, StringView, and mixed keys, different cardinalities and concurrency 
levels, per-query benchmark results, and peak memory usage.
   
   ### Describe alternatives you've considered
   
   - Keep the current behavior and rely on memory-pressure-triggered emission.
   - Skip Partial aggregation when its reduction is insufficient.
   
   These remain useful behaviors; proactive flushing could complement them by 
keeping aggregation effective with a smaller working set.
   
   ### Additional context
   
   Experiments so far:
   
   - #26100 isolates row and byte thresholds on main. In its same-binary 
comparison across all 43 partitioned ClickBench queries, the 2 MiB policy 
reduced the sum of per-query medians by approximately 4.6% and the geometric 
mean of query-time ratios by approximately 2.9%. Q33 and Q34 improved by 
approximately 9% each. Some queries regressed, so these are encouraging 
experimental results rather than evidence for a universal default.
   - #26116 explores target-batch flushing together with reservation, 
whole-batch emission, and progressive StringView payload blocks. Its combined 
results should not be attributed to flushing alone.
   
   I'd be happy to help investigate this further and follow up on the 
implementation and benchmarking.
   


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