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]
