github-actions[bot] commented on code in PR #68651:
URL: https://github.com/apache/doris/pull/68651#discussion_r4136186008
##########
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:
[P2] Track the per-bucket output key buffers. _output_bucket() resizes
output_keys in each of 256 bucket methods and retains their capacity until
shared-state teardown. At batch size 4096 with 32-byte fixed keys, buckets
containing at least one batch of groups retain about 32 MiB outside the Doris
allocator and operator memory counters. Use an allocator-aware reusable buffer,
ideally one per source task, and include its capacity in memory accounting.
##########
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:
[P1] Require the actual backend count before removing the exchange.
be_number_for_test is a settable session variable, so on a real multi-BE
cluster SET be_number_for_test=1 passes this gate. Scan ranges still go to
their actual backend workers, while bucketed aggregation merges only instances
within one BE. If the same group spans workers, the query returns duplicate
partial groups instead of one sum. Use the live cluster count for this
correctness check, independent of the test override.
##########
be/src/exec/operator/bucketed_aggregation_source_operator.cpp:
##########
@@ -0,0 +1,756 @@
+// 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.
+
+#include "exec/operator/bucketed_aggregation_source_operator.h"
+
+#include <memory>
+#include <string>
+
+#include "common/exception.h"
+#include "core/column/column_vector.h"
+#include "exec/common/hash_table/hash.h"
+#include "exec/common/util.hpp"
+#include "exec/operator/operator.h"
+#include "exprs/vectorized_agg_fn.h"
+#include "runtime/runtime_profile.h"
+#include "runtime/thread_context.h"
+
+namespace doris {
+
+// Helper to set/get null key data on hash tables that support it
(DataWithNullKey).
+// For hash tables without nullable key support (PHHashMap), these are no-ops.
+// This is needed because in nested std::visit lambdas, the outer hash table
type is already
+// resolved and doesn't depend on the inner template parameter, so `if
constexpr` inside the
+// inner lambda cannot suppress compilation of code that accesses
has_null_key_data() on the
+// outer (non-dependent) type.
+template <typename HashTable>
+constexpr bool has_nullable_key_v =
+
std::is_assignable_v<decltype(std::declval<HashTable&>().has_null_key_data()),
bool>;
+
+template <typename HashTable>
+void set_null_key_flag(HashTable& ht, bool val) {
+ if constexpr (has_nullable_key_v<HashTable>) {
+ ht.has_null_key_data() = val;
+ }
+}
+
+template <typename HashTable>
+bool get_null_key_flag(const HashTable& ht) {
+ if constexpr (has_nullable_key_v<HashTable>) {
+ return ht.has_null_key_data();
+ } else {
+ return false;
+ }
+}
+
+template <typename HashTable>
+AggregateDataPtr get_null_key_agg_data(HashTable& ht) {
+ if constexpr (has_nullable_key_v<HashTable>) {
+ return ht.template get_null_key_data<AggregateDataPtr>();
+ } else {
+ return nullptr;
+ }
+}
+
+template <typename HashTable>
+void set_null_key_agg_data(HashTable& ht, AggregateDataPtr val) {
+ if constexpr (has_nullable_key_v<HashTable>) {
+ ht.template get_null_key_data<AggregateDataPtr>() = val;
+ }
+}
+
+// Returns a REFERENCE to the null key's AggregateDataPtr slot.
+// Critical for simple_count merge: writing through a copy would lose the
update (Bug #30).
+template <typename HashTable>
+AggregateDataPtr& get_null_key_agg_data_ref(HashTable& ht) {
+ static_assert(has_nullable_key_v<HashTable>,
+ "get_null_key_agg_data_ref requires a nullable hash table");
+ return ht.template get_null_key_data<AggregateDataPtr>();
+}
+
+// Helper for emplace that works with PHHashMap (3-arg).
+template <typename HashTable, typename Key>
+auto hash_table_emplace(HashTable& ht, const Key& key, typename
HashTable::LookupResult& it,
+ bool& inserted) -> decltype(ht.emplace(key, it,
inserted), void()) {
+ ht.emplace(key, it, inserted);
+}
+
+/// Merge src aggregate state into dst_ref (a reference to the mapped slot).
+/// For simple_count, adds UInt64 counters directly via the reference.
+/// For regular aggregates, calls merge() on each function then destroys src
state.
+/// After return, src is consumed and must not be used.
+static void merge_agg_states(AggregateDataPtr& dst_ref, AggregateDataPtr src,
bool use_simple_count,
+ const std::vector<AggFnEvaluator*>& evaluators,
const Sizes& offsets,
+ Arena& arena) {
+ if (use_simple_count) {
+ // simple_count: mapped slots hold UInt64 counters. MUST use reference
+ // to write back correctly.
+ reinterpret_cast<UInt64&>(dst_ref) += reinterpret_cast<UInt64>(src);
+ } else {
+ const size_t num_fns = evaluators.size();
+ for (size_t i = 0; i < num_fns; ++i) {
+ evaluators[i]->function()->merge(dst_ref + offsets[i], src +
offsets[i], arena);
+ }
+ for (size_t i = 0; i < num_fns; ++i) {
+ evaluators[i]->function()->destroy(src + offsets[i]);
+ }
+ }
+}
+
+/// Merge a source null key into a destination null key slot. Handles three
cases:
+/// 1. Dst has no null key yet: move src's null key to dst (no merge needed).
+/// 2. Dst already has a null key: merge src into dst using merge_agg_states.
+/// 3. Src has no null key: no-op.
+/// After merge, clears the src null key slot.
+template <typename HashTable>
+static void merge_null_key(HashTable& dst_data, HashTable& src_data, bool
use_simple_count,
+ const std::vector<AggFnEvaluator*>& evaluators,
const Sizes& offsets,
+ Arena& arena) {
+ if constexpr (has_nullable_key_v<HashTable>) {
+ if (!get_null_key_flag(src_data)) {
+ return;
+ }
+ auto src_null = get_null_key_agg_data(src_data);
+ if (!src_null) {
+ return;
+ }
+ if (!get_null_key_flag(dst_data)) {
+ // Dst has no null key yet — move src's null key to dst.
+ set_null_key_flag(dst_data, true);
+ set_null_key_agg_data(dst_data, src_null);
+ } else {
+ // Both have null keys — merge src into dst.
+ auto& dst_null_ref = get_null_key_agg_data_ref(dst_data);
+ merge_agg_states(dst_null_ref, src_null, use_simple_count,
evaluators, offsets, arena);
+ }
+ set_null_key_agg_data(src_data, nullptr);
+ set_null_key_flag(src_data, false);
+ }
+}
+
+BucketedAggLocalState::BucketedAggLocalState(RuntimeState* state,
OperatorXBase* parent)
+ : Base(state, parent) {}
+
+Status BucketedAggLocalState::init(RuntimeState* state, LocalStateInfo& info) {
+ RETURN_IF_ERROR(Base::init(state, info));
+ SCOPED_TIMER(exec_time_counter());
+ SCOPED_TIMER(_init_timer);
+
+ _task_idx = info.task_idx;
+
+ _get_results_timer = ADD_TIMER(custom_profile(), "GetResultsTime");
+ _hash_table_iterate_timer = ADD_TIMER(custom_profile(),
"HashTableIterateTime");
+ _insert_keys_to_column_timer = ADD_TIMER(custom_profile(),
"InsertKeysToColumnTime");
+ _insert_values_to_column_timer = ADD_TIMER(custom_profile(),
"InsertValuesToColumnTime");
+ _merge_timer = ADD_TIMER(custom_profile(), "MergeTime");
+
+ return Status::OK();
+}
+
+Status BucketedAggLocalState::close(RuntimeState* state) {
+ SCOPED_TIMER(exec_time_counter());
+ SCOPED_TIMER(_close_timer);
+ if (_closed) {
+ return Status::OK();
+ }
+
+ // Release any held per-bucket CAS lock. This can happen when the source
+ // is closed prematurely (e.g., LIMIT reached via reached_limit() while
+ // we were mid-output on a bucket). Without this, the other source instance
+ // would spin forever trying to acquire this bucket's lock.
+ if (_current_output_bucket >= 0) {
+ auto& bs = _shared_state->bucket_states[_current_output_bucket];
+ bs.output_done.store(true, std::memory_order_release);
+ bs.merge_in_progress.store(false, std::memory_order_release);
+ _current_output_bucket = -1;
+ _shared_state->state_generation.fetch_add(1,
std::memory_order_release);
+ _wake_up_other_sources();
+ }
+
+ return Base::close(state);
+}
+
+void BucketedAggLocalState::_make_nullable_output_key(Block* block) {
+ if (block->rows() != 0) {
+ for (auto cid : _shared_state->make_nullable_keys) {
+ block->get_by_position(cid).column =
make_nullable(block->get_by_position(cid).column);
+ block->get_by_position(cid).type =
make_nullable(block->get_by_position(cid).type);
+ }
+ }
+}
+
+void BucketedAggLocalState::_wake_up_other_sources() {
+ auto& shared_state = *_shared_state;
+ for (int i = 0; i < static_cast<int>(shared_state.source_deps.size());
++i) {
+ shared_state.source_deps[i]->set_ready();
+ }
+}
+
+int BucketedAggLocalState::_merge_bucket(int bucket, int merge_target) {
+ SCOPED_TIMER(_merge_timer);
+ auto& shared_state = *_shared_state;
+ auto& bs = shared_state.bucket_states[bucket];
+ // Other source instances may merge other buckets at the same time, so
aggregate
+ // function merges must allocate from this source instance's own arena.
+ DCHECK_LT(_task_idx, shared_state.source_merge_arenas.size());
+ auto& merge_arena = *shared_state.source_merge_arenas[_task_idx];
+
+ // Merge target's bucket is the destination.
+ auto& dst_agg_data =
*shared_state.per_instance_data[merge_target].bucket_agg_data[bucket];
+ int merged_count = 0;
+
+ std::visit(
+ Overload {
+ [&](std::monostate& arg) -> void {
+ // uninited — no data to merge
+ },
+ [&](auto& dst_method) -> void {
+ using AggMethodType =
std::decay_t<decltype(dst_method)>;
+ auto& dst_data = *dst_method.hash_table;
+
+ // Merge all finished sink instances (except
merge_target itself)
+ // into the merge target's bucket.
+ for (int inst_idx = 0; inst_idx <
shared_state.num_sink_instances;
+ ++inst_idx) {
+ if (inst_idx == merge_target) {
+ continue;
+ }
+ // Skip instances already merged for this bucket.
+ if (bs.merged_instances[inst_idx]) {
+ continue;
+ }
+ // Only merge sinks that have finished.
+ if (!shared_state.sink_finished[inst_idx].load(
+ std::memory_order_acquire)) {
+ continue;
+ }
+
+ auto& src_inst =
shared_state.per_instance_data[inst_idx];
+ auto& src_agg_data =
*src_inst.bucket_agg_data[bucket];
+
+ std::visit(
+ Overload {
+ [&](std::monostate& arg) -> void {
+ // Mark as merged even if
monostate (no data).
+ bs.merged_instances[inst_idx]
= true;
+ },
+ [&](auto& src_method) -> void {
+ using SrcMethodType =
+
std::decay_t<decltype(src_method)>;
+ if constexpr
(std::is_same_v<SrcMethodType,
+
AggMethodType>) {
+ auto& src_data =
*src_method.hash_table;
+
+ ++merged_count;
+
+ // Direct merge: iterate
source hash table
+ // entries, emplace into
destination, and null
+ // out source entries in
one pass. This avoids
+ // allocating intermediate
vectors (keys,
+ // mappeds, hashes) and
eliminates the separate
+ // null-out traversal.
+ const bool
use_simple_count =
+
shared_state.use_simple_count;
+
src_data.for_each([&](const auto& key,
+
auto& mapped) {
+ if (!mapped) {
+ return;
+ }
+ auto src_mapped =
mapped;
+ mapped = nullptr;
Review Comment:
[P2] Keep the source state reachable until merge succeeds. This clears the
source mapped slot before destination emplace or aggregate merge(), both of
which can allocate and throw. For a COLLECT_SET group, an allocation failure
during its flat_hash_set merge unwinds past src_mapped; shared-state cleanup
then skips the null slot and never destroys the aggregate-owned set. Delay
clearing the slot until successful transfer/merge, with exception cleanup that
avoids double destruction.
##########
be/src/exec/pipeline/pipeline_fragment_context.cpp:
##########
@@ -1485,6 +1487,53 @@ Status
PipelineFragmentContext::_create_operator(ObjectPool* pool, const TPlanNo
}
break;
}
+ case TPlanNodeType::BUCKETED_AGGREGATION_NODE: {
Review Comment:
[P2] Preserve the spill path for spill-enabled grouped queries. This case
always builds BucketedAggSinkOperatorX and BucketedAggSourceOperatorX, while
the adjacent keyed aggregation case uses PartitionedAgg when enable_spill is
true. FE bucketed eligibility has no spill gate, so a single-BE,
high-cardinality GROUP BY can now hit its memory limit instead of spilling.
Fall back to regular aggregation when spilling is enabled until bucketed
aggregation implements revocation and spill.
##########
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:
[P2] Preserve evaluation of nonnullable COUNT arguments before using inline
count. is_simple_count() also applies to COUNT(assert_true(v > 0, 'bad')): on a
row with v <= 0, this branch enables the inline path, and the changed regular,
streaming, and bucketed sinks skip AggFnEvaluator, so the query returns a count
instead of raising the assertion error. Restrict the optimization to COUNT(*)
or evaluate nontrivial arguments before incrementing.
##########
regression-test/suites/mv_p0/ut/testBucketedAggSyncMV/testBucketedAggSyncMV.groovy:
##########
@@ -0,0 +1,100 @@
+// 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 ("testBucketedAggSyncMV") {
+ // Regression test: when enable_bucketed_hash_agg=true on a single-BE
cluster,
+ // the PhysicalBucketedHashAggregate plan for the base table with
multi_distinct_count
+ // could appear artificially cheap (bypasses one-phase agg ban, requests
ANY properties),
+ // beating the sync MV path on cost. The fix bans bucketed agg for
MultiDistinction
+ // functions in SplitAggWithoutDistinct.implementBucketedPhase().
+
+ sql "set pre_materialized_view_rewrite_strategy = TRY_IN_RBO"
+ sql """set enable_nereids_planner=true;"""
+ sql "set disable_nereids_rules='DISTINCT_AGGREGATE_SPLIT';"
+ sql "set enable_bucketed_hash_agg = true;"
+
+ sql """ DROP TABLE IF EXISTS bucketed_agg_mv_test; """
+
+ sql """
+ CREATE TABLE bucketed_agg_mv_test (
+ time_col date NOT NULL,
+ advertiser varchar(10),
+ dt date NOT NULL,
+ channel varchar(10),
+ user_id int)
+ DUPLICATE KEY(`time_col`, `advertiser`)
+ PARTITION BY RANGE (dt)(
+ FROM ("2024-07-01") TO ("2024-07-05") INTERVAL 1 DAY)
+ DISTRIBUTED BY hash(time_col) BUCKETS 3
+ PROPERTIES('replication_num' = '1');
+ """
+
+ sql """insert into bucketed_agg_mv_test
values("2024-07-01",'a',"2024-07-01",'x',1);"""
+ sql """insert into bucketed_agg_mv_test
values("2024-07-01",'a',"2024-07-01",'x',1);"""
+ sql """insert into bucketed_agg_mv_test
values("2024-07-02",'a',"2024-07-02",'y',2);"""
+ sql """insert into bucketed_agg_mv_test
values("2024-07-02",'b',"2024-07-02",'x',1);"""
+ sql """insert into bucketed_agg_mv_test
values("2024-07-03",'b',"2024-07-03",'y',3);"""
+
+ createMV("""
+ CREATE MATERIALIZED VIEW bucketed_agg_uv_mv AS
+ SELECT advertiser AS a1,
+ channel AS a2,
+ dt AS a3,
+ bitmap_union(to_bitmap(user_id)) AS a4
+ FROM bucketed_agg_mv_test
+ GROUP BY advertiser,
+ channel,
+ dt;
+ """)
+
+ sql """insert into bucketed_agg_mv_test
values("2024-07-03",'b',"2024-07-03",'y',4);"""
+
+ sql "analyze table bucketed_agg_mv_test with sync;"
+ sql """alter table bucketed_agg_mv_test modify column time_col set stats
('row_count'='6');"""
+
+ // Core assertion: with enable_bucketed_hash_agg=true, the sync MV must
still be chosen.
+ // Before the fix, PhysicalBucketedHashAggregate(multi_distinct_count)
would win on cost.
+ mv_rewrite_success(
Review Comment:
[P2] Make the bucketed alternative eligible before asserting the MV wins.
This fixture analyzes six rows and pins row_count=6, but leaves
bucketed_agg_min_input_rows at its 100000 default. The new regulator rejects
the one-phase bucketed candidate at that gate, so mv_rewrite_success passes
even if bucketed aggregation would beat the MV when eligible. Set suitable
bucketed thresholds for this fixture and assert a comparable base-table query
actually takes the bucketed plan as a positive control.
--
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]