github-actions[bot] commented on code in PR #68651:
URL: https://github.com/apache/doris/pull/68651#discussion_r4151771373


##########
fe/fe-core/src/main/java/org/apache/doris/nereids/util/AggregateUtils.java:
##########
@@ -161,4 +189,346 @@ public static Set<NamedExpression> 
getDistinctNamedExpr(LogicalAggregate<? exten
                 .map(NamedExpression.class::cast)
                 .collect(ImmutableSet.toImmutableSet());
     }
+
+    /**
+     * Check if order keys are identical to group-by keys (1-1 mapping, same 
order).
+     * Shared utility used by both PushTopnToAgg and SplitAggWithoutDistinct.
+     */
+    public static boolean isOrderKeysMatchGroupKeys(List<OrderKey> orderKeys,
+            List<Expression> groupByKeys) {
+        if (orderKeys.size() != groupByKeys.size()) {
+            return false;
+        }
+        for (int i = 0; i < groupByKeys.size(); i++) {
+            if (!groupByKeys.get(i).equals(orderKeys.get(i).getExpr())) {
+                return false;
+            }
+        }
+        return true;
+    }
+
+    /**
+     * Check the basic environmental conditions for bucketed hash aggregation.
+     * This is the environment part of the shared eligibility gate; the 
physical
+     * plan shape part is {@link 
#isBucketedHashAggFusible(PhysicalHashAggregate)},
+     * {@link #isBucketedHashAggFusible(PhysicalHashAggregate, 
DistributionSpec)} and
+     * {@link #isBucketedHashAggFusible(GroupExpression, PhysicalProperties)},
+     * which ChildrenPropertiesRegulator (to allow the 
one-phase-GLOBAL+distribute
+     * pattern), ChildOutputPropertyDeriver, CostModel (for the cost discount) 
and
+     * PhysicalPlanTranslator (for fusion into BucketedAggregationNode) all 
use.
+     *
+     * @return true if the session variable is enabled, there is exactly one 
alive BE,
+     *         spill and the query cache are disabled, no smooth upgrade is in 
progress,
+     *         the aggregate has GROUP BY keys and contains no user-defined 
aggregate function.
+     */
+    public static boolean isBucketedHashAggEnabled(Aggregate<? extends Plan> 
aggregate) {
+        ConnectContext ctx = ConnectContext.get();
+        if (ctx == null) {
+            return false;
+        }
+        if (!ctx.getSessionVariable().enableBucketedHashAgg) {
+            return false;
+        }
+        // Must have GROUP BY keys (without-key aggregation not supported)
+        if (aggregate.getGroupByExpressions().isEmpty()) {
+            return false;
+        }
+        // Bucketed agg has no spill support. Keep the regular (spillable) 
aggregation
+        // when spill is enabled, otherwise a high-cardinality GROUP BY could 
hit the
+        // memory limit instead of spilling.
+        if (ctx.getSessionVariable().enableSpill || 
ctx.getSessionVariable().enableForceSpill) {
+            return false;
+        }
+        // The query cache is built on the LOCAL AggregationNode above the 
scan, and neither
+        // the FE normalizer nor the BE cache operators know 
BucketedAggregationNode. Keep the
+        // regular aggregation when the query cache is enabled, otherwise the 
query would
+        // silently lose its cache point.
+        if (ctx.getSessionVariable().getEnableQueryCache()) {
+            return false;
+        }
+        // be_number_for_test can only disable bucketed agg (to test the 
multi-BE plan),
+        // never bypass the single-BE gate below.
+        int beNumberForTest = ctx.getSessionVariable().getBeNumberForTest();
+        if (beNumberForTest > 0 && beNumberForTest != 1) {
+            return false;
+        }
+        // Correctness gate: single-BE only (cross-BE in-memory merge is 
impossible).
+        if (getBucketedHashAggBackend(ctx) == null) {
+            return false;
+        }
+        // Bucketed agg merges the live states built by different sink 
instances
+        // directly, without serializing them. Java / Python UDAFs can only 
merge a
+        // state that was deserialized by the merging evaluator (the Java UDAF
+        // executor place and the Python UDAF serialized buffer are only set 
up on
+        // that path), so they must stay on the regular aggregation path.
+        if 
(aggregate.getAggregateFunctions().stream().anyMatch(Udf.class::isInstance)) {
+            return false;
+        }
+        return true;
+    }
+
+    /**
+     * Returns the backend that a bucketed hash aggregation of this query 
would run on:
+     * the only alive backend of the current cluster. Returns null when the 
number of
+     * alive backends is not exactly one, or when that backend is a 
smooth-upgrade source
+     * (its old BE process does not know BUCKETED_AGGREGATION_NODE).
+     * <p>
+     * Scan ranges always go to the real alive backends, so they are counted 
directly
+     * (getBackendsNumber() would return be_number_for_test), and zero 
backends must not
+     * be clamped to one. PhysicalPlanTranslator reads the backend again when 
it fuses
+     * an aggregate and pins the scan of the fused fragment to it, so a 
backend that
+     * becomes alive after this check cannot receive tablets of that fragment.
+     */
+    public static Backend getBucketedHashAggBackend(ConnectContext ctx) {
+        SystemInfoService clusterInfo = ctx.getEnv().getClusterInfo();
+        List<Long> aliveBackendIds = 
clusterInfo.getAllBackendByCurrentCluster(true);
+        if (aliveBackendIds.size() != 1) {
+            return null;
+        }
+        Backend be = clusterInfo.getBackend(aliveBackendIds.get(0));
+        if (be == null || be.isSmoothUpgradeSrc()) {
+            return null;
+        }
+        return be;
+    }
+
+    /**
+     * Check whether a physical hash aggregate has the shape that 
PhysicalPlanTranslator
+     * fuses into BucketedAggregationNode, as far as the aggregate itself is 
concerned:
+     * the environment gate above passes, it is the one-phase GLOBAL 
INPUT_TO_RESULT
+     * aggregate, none of its functions produces a partial buffer, all of them 
support
+     * two-phase execution and no TopN was pushed into it. The translator 
additionally
+     * requires a distribute child on exactly the GROUP BY keys over a unary 
single-scan
+     * pipeline, which is not visible here.
+     * <p>
+     * The regulator, the output property deriver and the cost model must use 
this
+     * gate rather than {@link #isBucketedHashAggEnabled(Aggregate)} alone: an 
aggregate
+     * that passes the environment gate but not the shape gate (for example 
the GLOBAL
+     * INPUT_TO_RESULT dedup aggregate of a mixed DISTINCT / non-DISTINCT 
query, whose
+     * non-distinct functions run in INPUT_TO_BUFFER mode) is translated into 
a regular
+     * AggregationNode that keeps the exchange, so it must not receive the 
bucketed
+     * cost discount or the one-phase-with-distribute exemption. Callers that 
know the
+     * distribution of the aggregate's child use
+     * {@link #isBucketedHashAggFusible(PhysicalHashAggregate, 
DistributionSpec)}, and
+     * callers that work on the memo, where the child's subtree is visible as 
well, use
+     * {@link #isBucketedHashAggFusible(GroupExpression, PhysicalProperties)}.
+     */
+    public static boolean isBucketedHashAggFusible(PhysicalHashAggregate<? 
extends Plan> aggregate) {
+        // Must be one-phase: GLOBAL + INPUT_TO_RESULT. Checked before the 
environment
+        // gate, which asks the cluster for its alive backends, because every 
aggregate
+        // alternative in the memo passes through here.
+        if (aggregate.getAggPhase() != AggPhase.GLOBAL
+                || aggregate.getAggMode() != AggMode.INPUT_TO_RESULT) {
+            return false;
+        }
+        if (!isBucketedHashAggEnabled(aggregate)) {
+            return false;
+        }
+        // BucketedAggregationNode always finalizes into the output tuple slot
+        // types (need_finalize=true, isPartial=false), so fusing an aggregate
+        // whose functions produce buffers (partial) would fail the BE
+        // result-type check: the slot type of a buffer-producing function is
+        // Varchar while the function's final return type (e.g. DOUBLE for
+        // stddev) is what insert_result_into writes. The one-phase GLOBAL
+        // dedup aggregate of a 3-phase DISTINCT plan has exactly this shape —
+        // the node itself is INPUT_TO_RESULT but its non-distinct functions
+        // run in INPUT_TO_BUFFER mode — and must stay on the regular
+        // AggregationNode path, which serializes when isPartial.
+        if (containsPartialAggFunction(aggregate)) {
+            return false;
+        }
+        // Exclude one-phase-only aggregates (e.g. GROUP_CONCAT with ORDER BY).
+        // BucketedAggregationNode has no sort-info field, so fusing would drop
+        // the aggregate ORDER BY contract. Only aggregates supporting 
two-phase
+        // execution can be safely fused.
+        if (!supportsTwoPhaseAgg(aggregate)) {
+            return false;
+        }
+        // BucketedAggregationNode does not support sortByGroupKey 
(PushTopnToAgg
+        // optimization). Regular AggregationNode fills sort info; fusing 
would drop it.
+        return aggregate.getTopnPushInfo() == null;

Review Comment:
   [P2] Keep the bucketed cost decision valid after TopN pushdown. When the 
optimizer selects `LOCAL_SORT TopN(order by k LIMIT 10) -> GlobalAgg(k, 
COUNT(*)) -> Distribute(HASH(k)) -> OlapScan` on a qualifying single BE, memo 
regulation admits the one-phase aggregate and `CostCalculator` halves its row 
cost while `getTopnPushInfo()` is null. The later `PushTopnToAgg` postprocessor 
sets that hint; this check then rejects fusion during translation, leaving the 
discounted plan with a regular aggregate and a raw-row exchange. Account for 
this TopN interaction before granting the one-phase exemption/discount, or make 
the fused node honor the pushed TopN; add a planner test for this shape.



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