mrhhsg commented on code in PR #68651:
URL: https://github.com/apache/doris/pull/68651#discussion_r4139546877
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/util/AggregateUtils.java:
##########
@@ -161,4 +166,73 @@ 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,
+ * 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;
+ }
+ // Correctness gate: single-BE only (cross-BE in-memory merge is
impossible).
+ // Use be_number_for_test first (set by regression tests), fall back
to real cluster count.
+ // Note: do not clamp to 1 — with zero backends bucketed agg must not
be enabled.
+ int beNumber = ctx.getSessionVariable().getBeNumberForTest();
Review Comment:
Fixed in 33bbe3482ea. `isBucketedHashAggEnabled` now counts the real alive
backends of the current cluster with `getAllBackendByCurrentCluster(true)` (not
`getBackendsNumber()`, which returns `be_number_for_test` when it is set) and
requires exactly one. `be_number_for_test` can now only disable the fusion (`>
0 && != 1`), so tests can still check the multi-BE plan but cannot enable
bucketed aggregation on a multi-BE cluster. New FE UT
`BucketedAggregateMultiBackendTest` starts 2 backends, sets
`be_number_for_test=1`, and asserts that no `BucketedAggregationNode` is
planned. It fails with the old gate.
##########
be/src/exec/operator/aggregation_sink_operator.cpp:
##########
@@ -156,6 +157,30 @@ Status AggSinkLocalState::open(RuntimeState* state) {
RETURN_IF_ERROR(_create_agg_status(_agg_data->without_key));
_shared_state->agg_data_created_without_key = true;
}
+
+ // Determine whether to use simple count aggregation.
+ // For queries like: SELECT xxx, count(*) / count(not_null_column) FROM
table GROUP BY xxx,
+ // count(*) / count(not_null_column) can store a uint64 counter directly
in the hash table,
+ // instead of storing the full aggregate state, saving memory and
computation overhead.
+ // Requirements:
+ // 0. The aggregation has a GROUP BY clause.
+ // 1. There is exactly one count aggregate function.
+ // 2. No limit optimization is applied.
+ // 3. Spill is not enabled (the spill path accesses
aggregate_data_container, which is empty in inline count mode).
+ // Supports update / merge / finalize / serialize phases, since count's
serialization format is UInt64 itself.
+
+ if (!Base::_shared_state->probe_expr_ctxs.empty() /* has GROUP BY */
+ && (p._aggregate_evaluators.size() == 1 &&
+ p._aggregate_evaluators[0]->function()->is_simple_count()) /* only
one count(*) */
Review Comment:
Fixed in 33bbe3482ea. Added `AggFnEvaluator::is_simple_count()`: it is true
only when the function is a simple count and every argument root is a slot ref
or a literal. The regular, streaming and bucketed sinks now use it, so
`COUNT(assert_true(...))` goes through `AggFnEvaluator` and raises the error.
`COUNT(*)`, `COUNT(non-null column)` and the merge phases still use the inline
path. Covered by BE UT `AggFnEvaluatorTest.test_is_simple_count` and by
regression `bucketed_hash_agg` Test 8, which checks that `count(assert_true(val
< 100, ...)) GROUP BY grp` raises the error with and without the bucketed plan.
##########
be/src/exec/common/hash_table/hash_map_context.h:
##########
@@ -54,6 +54,11 @@ struct MethodBaseInner {
Arena arena;
DorisVector<size_t> hash_values;
+ /// Reusable buffer for source-side output iteration to avoid per-batch
+ /// heap allocation of std::vector<Key>. Callers use resize() + direct
+ /// element assignment, so the capacity is retained across batches.
+ std::vector<Key> output_keys;
Review Comment:
Fixed in 33bbe3482ea. `_output_bucket` now releases the bucket method's
`output_keys` (swap with an empty vector) as soon as the bucket's iterator
reaches the end. The check comes before any `resize`, so the final call that
confirms the end allocates nothing. At most the bucket each source instance is
currently outputting keeps a batch-sized buffer, instead of up to 256 retained
buffers until shared-state teardown.
--
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]