github-actions[bot] commented on code in PR #68651: URL: https://github.com/apache/doris/pull/68651#discussion_r4139772815
########## be/src/exec/operator/bucketed_aggregation_sink_operator.h: ########## @@ -0,0 +1,147 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#pragma once + +#include <stdint.h> + +#include "exec/operator/operator.h" +#include "runtime/exec_env.h" +#include "runtime/runtime_profile.h" + +namespace doris { + +class BucketedAggSinkOperatorX; + +/// Sink-side local state for bucketed hash aggregation. +/// Each pipeline instance builds 256 per-bucket hash tables (two-level hash table). +/// No locking: each instance writes to per_instance_data[_instance_idx]. +class BucketedAggSinkLocalState : public PipelineXSinkLocalState<BucketedAggSharedState> { +public: + ENABLE_FACTORY_CREATOR(BucketedAggSinkLocalState); + using Base = PipelineXSinkLocalState<BucketedAggSharedState>; + BucketedAggSinkLocalState(DataSinkOperatorXBase* parent, RuntimeState* state); + ~BucketedAggSinkLocalState() override = default; + + Status init(RuntimeState* state, LocalSinkStateInfo& info) override; + Status open(RuntimeState* state) override; Review Comment: [P2] Mark the bucketed sink blockable for blocking aggregate functions. `BucketedAggSinkLocalState` inherits `is_blockable() == false`, while the regular aggregate sink checks its evaluator. An eligible grouped `AI_AGG` query can issue a synchronous HTTP request from `add()` after its context fills, yet `PipelineTask::is_blockable()` will route this task to the CPU scheduler. Implement the evaluator check here (and account for source-side merge/finalization) so the request cannot occupy a CPU worker for its timeout. ########## fe/fe-core/src/main/java/org/apache/doris/nereids/util/AggregateUtils.java: ########## @@ -161,4 +166,83 @@ 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 shared eligibility gate used by ChildrenPropertiesRegulator + * (to allow the one-phase-GLOBAL+distribute pattern), CostModel (for cost + * discount), and PhysicalPlanTranslator (for fusion into BucketedAggregationNode). + * + * @return true if the session variable is enabled, there is exactly one alive BE, + * spill is 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; + } + // 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). + // Scan ranges always go to the real alive backends, so count them directly + // (getBackendsNumber() would return be_number_for_test). + // Note: do not clamp to 1 — with zero backends bucketed agg must not be enabled. + SystemInfoService clusterInfo = ctx.getEnv().getClusterInfo(); + List<Long> aliveBackendIds = clusterInfo.getAllBackendByCurrentCluster(true); + if (aliveBackendIds.size() != 1) { + return false; + } + // Smooth upgrade safety net: old BE processes do not recognize + // BUCKETED_AGGREGATION_NODE plan node type + for (Long beId : aliveBackendIds) { Review Comment: [P1] Gate the new plan node on the backend's actual capability. This check only excludes `isSmoothUpgradeSrc`, which is set for cloud smooth upgrade and remains false for a shared-nothing BE. After an FE-first upgrade with the sole BE still on the old version, the default-enabled bucketed path sends `BUCKETED_AGGREGATION_NODE` to that BE; its dispatcher has no case for the new enum, so eligible grouped queries fail. Check a BE capability or minimum version before enabling this path. ########## fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslator.java: ########## @@ -3216,6 +3202,351 @@ private PlanFragment connectJoinNode(HashJoinNode hashJoinNode, PlanFragment lef return leftFragment; } + /** + * Check whether the one-phase GLOBAL hash aggregate can be fused with its + * distribute child into a BucketedAggregationNode. This eliminates exchange + * overhead on single-BE deployments by using in-memory per-bucket merging. + */ + private boolean shouldUseBucketedFusion(PhysicalHashAggregate<? extends Plan> aggregate, + PlanTranslatorContext context) { + // Shared eligibility: session var, single-BE, GROUP BY, smooth upgrade, no UDAF + if (!AggregateUtils.isBucketedHashAggEnabled(aggregate)) { + return false; + } + // Must be one-phase: GLOBAL + INPUT_TO_RESULT + if (aggregate.getAggPhase() != AggPhase.GLOBAL + || aggregate.getAggMode() != AggMode.INPUT_TO_RESULT) { + 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. + if (aggregate.getTopnPushInfo() != null) { + return false; + } + // Child must be PhysicalDistribute with hash distribution matching group keys + Plan child = aggregate.child(0); + if (!(child instanceof PhysicalDistribute)) { + return false; + } + // Bucketed fusion bypasses the distribute/exchange and builds directly on the + // child fragment. When the child subtree contains a CTE consumer (materialized + // multicast CTE), the child fragment is the MultiCastPlanFragment; a parent + // distribute would then treat the aggregate output slots as consumer slots and + // fail with "Required producer slot ... doesn't exist". Fall back to the + // regular one-phase path (which keeps the exchange) for such plans. + if (containsCTEConsumer(child)) { + return false; + } + // The distribute's child subtree must be a unary pipeline over exactly one + // olap scan. Fusing an aggregate whose input contains a join / set-op / CTE + // subtree would leave multiple olap scans in a single fragment (rejected by + // UnassignedJobBuilder: "Not supported multiple scan multiple OlapTable but + // not contains colocate join or bucket shuffle join"), and fusing over a + // nested aggregate would break the bucket alignment between stages. + if (!isSingleOlapScanPipeline(aggregate.child(0).child(0))) { + return false; + } + // The parent is a fragment-merging node (join / set-op) that consumes this + // fragment without an exchange boundary: fusing removes the exchange that + // keeps the scan in its own fragment, so multiple scans would end up in the + // same fragment and the scan-assignment would fail. Only fuse when the + // parent chain keeps an exchange boundary (e.g. a top-level aggregate). + if (context.isInFragmentMergeChild()) { + return false; + } + DistributionSpec distSpec = ((PhysicalDistribute<?>) child).getDistributionSpec(); + if (!(distSpec instanceof DistributionSpecHash)) { + return false; + } + List<ExprId> distKeys = ((DistributionSpecHash) distSpec).getOrderedShuffledColumns(); + List<ExprId> groupByKeys = aggregate.getGroupByExpressions().stream() + .filter(SlotReference.class::isInstance) + .map(SlotReference.class::cast) + .map(SlotReference::getExprId) + .collect(Collectors.toList()); + return distKeys.equals(groupByKeys); + } + + /** Returns true if the plan subtree contains a physical CTE consumer. */ + private boolean containsCTEConsumer(Plan plan) { + if (plan instanceof PhysicalCTEConsumer) { + return true; + } + for (Plan child : plan.children()) { + if (containsCTEConsumer(child)) { + return true; + } + } + return false; + } + + /** + * Returns true if the plan subtree is a unary pipeline over exactly one olap + * scan, i.e. it translates into a single-scan fragment that bucketed fusion + * can safely build upon. Subtrees containing fragment-merging or + * distribution-changing nodes (join / set-op / CTE / nested aggregate / + * storage-layer aggregate) are rejected. + */ + private boolean isSingleOlapScanPipeline(Plan plan) { + if (plan instanceof PhysicalOlapScan) { + return true; + } + if (plan instanceof PhysicalHashJoin + || plan instanceof PhysicalNestedLoopJoin + || plan instanceof PhysicalSetOperation + || plan instanceof PhysicalCTEConsumer + || plan instanceof PhysicalCTEAnchor + || plan instanceof PhysicalHashAggregate + || plan instanceof PhysicalStorageLayerAggregate) { + return false; + } + if (plan.children().size() == 1) { + return isSingleOlapScanPipeline(plan.child(0)); + } + return false; + } + + /** + * Check whether all aggregate functions in this physical hash aggregate + * support two-phase execution. One-phase-only aggregates (e.g. GROUP_CONCAT + * with ORDER BY) cannot be bucketed because BucketedAggregationNode does not + * carry sort-info metadata (aggSortInfos); fusing them would drop the + * aggregate ORDER BY contract and produce unordered results. + */ + private boolean supportsTwoPhaseAgg(PhysicalHashAggregate<? extends Plan> aggregate) { + for (NamedExpression o : aggregate.getOutputExpressions()) { + AtomicBoolean foundOnePhaseOnly = new AtomicBoolean(false); + o.foreach(c -> { + if (c instanceof OrderExpression) { + // Any aggregate function with an internal ORDER BY + // (e.g. GROUP_CONCAT(... ORDER BY ...)) needs sort-info + // metadata, which BucketedAggregationNode does not carry. + foundOnePhaseOnly.set(true); + return true; + } + if (c instanceof AggregateExpression) { + AggregateFunction func = ((AggregateExpression) c).getFunction(); + if (!func.supportAggregatePhase(AggregatePhase.TWO)) { + foundOnePhaseOnly.set(true); + return true; + } + } + return false; + }); + if (foundOnePhaseOnly.get()) { + return false; + } + } + return true; + } + + /** + * Check whether the aggregate's output contains any buffer-producing + * (partial) aggregate function, i.e. an AggregateExpression in a mode with + * productAggregateBuffer=true. BucketedAggregationNode cannot carry such + * functions: it always finalizes into the output tuple slot types, while a + * buffer-producing function's slot type is the serialized Varchar type + * (AggregateExpression.getDataType) — writing the final result (e.g. DOUBLE + * for stddev) into that String column fails the BE result-type check. + */ + private boolean containsPartialAggFunction(PhysicalHashAggregate<? extends Plan> aggregate) { + for (NamedExpression o : aggregate.getOutputExpressions()) { + AtomicBoolean foundPartial = new AtomicBoolean(false); + o.foreach(c -> { + if (c instanceof AggregateExpression) { + if (((AggregateExpression) c).getAggregateParam().aggMode.productAggregateBuffer) { + foundPartial.set(true); + } + return true; + } + return false; + }); + if (foundPartial.get()) { + return true; + } + } + return false; + } + + /** + * Fuse a one-phase GLOBAL hash aggregate and its PhysicalDistribute child + * into a BucketedAggregationNode, skipping the exchange node entirely. + * Visits the distribute's child directly to keep everything in one fragment. + */ + private PlanFragment visitBucketedFusion( + PhysicalHashAggregate<? extends Plan> aggregate, + PlanTranslatorContext context) { + // Visit the distribute's direct child, bypassing the distribute entirely. + // This avoids creating an ExchangeNode that bucketed agg does not need. + Plan distributeChild = aggregate.child(0).child(0); + PlanFragment inputPlanFragment = distributeChild.accept(this, context); + + List<Expression> groupByExpressions = aggregate.getGroupByExpressions(); + List<NamedExpression> outputExpressions = aggregate.getOutputExpressions(); + + // 1. generate slot reference for each group expression + List<SlotReference> groupSlots = collectGroupBySlots(groupByExpressions, outputExpressions); + ArrayList<Expr> execGroupingExpressions = translateGroupByExprs(groupByExpressions, context); + + // 2. collect agg expressions and generate agg function to slot reference map. + // Reuse the shared helper from visitPhysicalHashAggregate; the bucketed + // path passes null for hasPartialOut (never partial, always needsFinalize). + Pair<List<Slot>, ArrayList<FunctionCallExpr>> aggResult = + collectAggFunctions(outputExpressions, null, context); + List<Slot> aggFunctionOutput = aggResult.first; + ArrayList<FunctionCallExpr> execAggregateFunctions = aggResult.second; + + // 3. generate output tuple + Pair<TupleDescriptor, List<Integer>> tupleAndIds = + buildAggOutputTuple(groupSlots, aggFunctionOutput, context); + TupleDescriptor outputTupleDesc = tupleAndIds.first; + List<Integer> aggFunOutputIds = tupleAndIds.second; + + // Bucketed agg uses AggPhase.FIRST (update semantics): raw input -> final result. + // Not partial — always needsFinalize. + AggregateInfo aggInfo = AggregateInfo.create(execGroupingExpressions, execAggregateFunctions, + aggFunOutputIds, false /* isPartial */, outputTupleDesc, + AggregateInfo.AggPhase.FIRST); + + BucketedAggregationNode bucketedAggNode = new BucketedAggregationNode( + context.nextPlanNodeId(), inputPlanFragment.getPlanRoot(), aggInfo, true); Review Comment: [P2] Preserve query-cache admission for cache-enabled grouped queries. The regular two-phase plan has a cacheable LOCAL `AggregationNode` over the scan, but this fused path installs a `BucketedAggregationNode` without that cache point. `QueryCacheNormalizer` recognizes only `AggregationNode`, and the BE creates cache operators only for that node type, so with `enable_query_cache=true` an eligible grouped OLAP query loses cache hits. Skip fusion when query cache is enabled until both FE normalization and BE cache operators support this node. ########## regression-test/suites/nereids_rules_p0/agg_strategy/cse_agg_distribute.groovy: ########## @@ -0,0 +1,69 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +suite("cse_agg_distribute") { + sql "SET enable_nereids_planner=true" + sql "SET enable_fallback_to_original_planner=false" + sql "SET runtime_filter_mode=OFF" + + sql "DROP TABLE IF EXISTS cse_agg_distribute_tbl" + sql """ + CREATE TABLE cse_agg_distribute_tbl ( + id int, + grp varchar(20), + a int, + b int + ) DUPLICATE KEY(id) + DISTRIBUTED BY HASH(id) BUCKETS 3 + PROPERTIES('replication_num' = '1') + """ + sql """ INSERT INTO cse_agg_distribute_tbl VALUES + (1, 'g1', 1, 2), + (2, 'g2', 3, 4), + (3, 'g1', 5, 6), + (4, 'g2', 7, 8), + (5, 'g1', 9, 10) + """ + + // SUM(a+b) and MAX(a+b) share the same argument, so the aggregate-argument + // CSE must extract "a+b" into a project node and make both functions + // reference the extracted slot, instead of re-evaluating a+b per function. + String query = "SELECT grp, SUM(a+b), MAX(a+b) FROM cse_agg_distribute_tbl GROUP BY grp" + + // --------------------------------------------------------------------- + // one-phase aggregate over a distribute (the aggregate is a join child, + // so the distribute is required by the join): the CSE project must be + // inserted below the distribute, keeping the distribution-key slots + // intact. Both aggregates must reference the extracted slot (4 + // occurrences: SUM/MAX of each side). + // --------------------------------------------------------------------- + sql "set agg_phase=1" + sql "set enable_bucketed_hash_agg=false" + String joinQuery = """ + SELECT t1.grp, t1.s, t1.m, t2.s2, t2.m2 FROM + (SELECT grp, SUM(a+b) s, MAX(a+b) m FROM cse_agg_distribute_tbl GROUP BY grp) t1 + JOIN (SELECT grp, SUM(a+b) s2, MAX(a+b) m2 FROM cse_agg_distribute_tbl GROUP BY grp) t2 + ON t1.grp = t2.grp + """ + explain { + sql("${joinQuery}") + contains("VEXCHANGE") + contains("VSELECT") + multiContains("cast(a as BIGINT) + cast(b as BIGINT))[#", 4) Review Comment: [P2] Fix the malformed `multiContains` assertion. The call closes after the first argument and leaves `[#", 4)` with mismatched delimiters, so Groovy cannot parse this new suite and none of its CSE checks run. Use `multiContains("cast(a as BIGINT) + cast(b as BIGINT)", 4)`. -- 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]
