mrhhsg opened a new pull request, #68752:
URL: https://github.com/apache/doris/pull/68752
### What problem does this PR solve?
Issue Number: None
Problem Summary:
Bucketed hash aggregation (fusing a one-phase GLOBAL aggregate with its
distribute into `BucketedAggregationNode` on a single-BE cluster, skipping the
exchange) has several follow-up issues:
1. **CSE double-projects over an existing project**: when the CSE
post-processor inserts a project below a one-phase aggregate's distribute
child, and that child is already a `PhysicalProject` (e.g. a projected consumer
of a materialized CTE), it used to stack a second project instead of merging.
Translating two stacked projects onto the same multicast sink then raised
`generate invalid plan`. Fixed by merging into the existing project
(`Project.canMergeChildProjections`/`mergeProjections`), or dropping the CSE
rewrite when they cannot be merged.
2. **No spill fallback**: `BucketedAggregationNode` has no
`revocable_mem_size`/`revoke_memory`, so a high-cardinality GROUP BY under
`enable_spill` could only grow memory until OOM instead of spilling. Fixed by
disabling bucketed agg whenever spill is enabled, and by making the BE source
operator account for merge-time memory growth (two new profile counters plus
`SCOPED_PEAK_MEM`) so the pipeline task's memory reservation reflects the real
cost of merging.
3. **Missing `is_blockable`**: neither the bucketed sink nor the source
reported `is_blockable()`, so a potentially blocking aggregate function (e.g. a
remote/AI aggregate) could occupy a non-blocking scheduler worker. Fixed by
having the sink report whether any of its aggregate functions is blockable, and
the source consult its paired sink.
4. **Query cache silently disabled**: the query cache point is the `LOCAL
AggregationNode` above the scan; neither the FE normalizer nor the BE cache
operators recognize `BucketedAggregationNode`. Fixed by disabling bucketed agg
whenever the query cache is enabled.
5. **Inconsistent fusion gating**: the regulator, the output-property
deriver, and the cost model each used a looser eligibility check than what the
translator actually fuses, so some shapes (a mixed DISTINCT dedup aggregate, a
nested aggregate, a projected CTE consumer, a parent-key-subset distribute)
were costed/planned as if they would be fused but then kept as a regular
aggregate by the translator — ending up strictly worse than the normal
two-phase plan. Fixed by centralizing the shape check into
`AggregateUtils.isBucketedHashAggFusible` (plain, child-distribution, and
memo-aware overloads) and having all four call sites use it. Also, a real
exchange (`PhysicalDistribute`) now clears the "inside a fragment-merging
parent" context, so an aggregate below a join/union's shuffle exchange can be
fused correctly.
6. **Merge failure orphans source state**:
`BucketedAggLocalState::_merge_bucket` used to null out the source hash-table
slot before the destination `emplace`/`merge_agg_states` call, so a thrown
exception during that call left the source state neither owned by the
destination nor visible for cleanup — a leak. Fixed by only clearing the source
slot after it is successfully handed off.
7. **`is_simple_count` swallows errors**: the inline-count fast path
(`is_simple_count()`) applied to `COUNT(...)` regardless of its argument,
skipping evaluation entirely — so `COUNT(assert_true(...))` or similar
expressions that must raise an error were silently skipped. Fixed by requiring
every argument to be a slot reference or a literal.
### Release note
Bucketed hash aggregation now has additional correctness and consistency
gates:
- It is disabled whenever `enable_spill`/`enable_force_spill` or the query
cache is enabled; the regular (spillable/cacheable) aggregation path is used
instead.
- The optimizer's decision to fuse an aggregate into
`BucketedAggregationNode` is now consistent across cost estimation,
physical-property derivation, and plan enforcement, so some plan shapes that
previously degraded to a regular aggregate shuffling raw rows will correctly
use two-phase aggregation instead.
- An aggregate below a join/union's shuffle exchange can now be fused,
avoiding a redundant exchange.
- A query whose `COUNT(...)` argument must raise an error (e.g.
`COUNT(assert_true(...))`) now does so under bucketed hash aggregation as well.
### Check List (For Author)
- Test:
- Unit Test: Added/updated BE UT (`agg_fn_evaluator_test.cpp`,
`agg_operator_test.cpp`) and FE UT (`BucketedAggregateTranslatorTest`,
`RecursiveUnionFragmentMergeContextTest`, new
`ProjectAggregateExpressionsForCseTest`); verified each new/changed assertion
fails on the pre-fix code and passes after the fix.
- Regression test: Added/updated
`nereids_rules_p0/agg_strategy/bucketed_hash_agg.groovy` and
`cse_agg_distribute.groovy`; added
`enable_spill=false`/`enable_force_spill=false` to several existing
bucketed-agg regression suites so CI's fuzzy spill setting does not make them
flaky. Also ran `nereids_tpch_p0/tpch` and `shape_check/clickbench` locally to
confirm the tightened cost/regulator gating does not change plan shapes outside
the targeted scenarios.
- Behavior changed: Yes — see Release note above.
- Does this need documentation: No
--
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]