github-actions[bot] commented on code in PR #68651:
URL: https://github.com/apache/doris/pull/68651#discussion_r4152305677
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/cost/CostCalculator.java:
##########
@@ -43,6 +45,12 @@ public static Cost calculateCost(ConnectContext
connectContext, GroupExpression
&& childrenProperties.get(1).getDistributionSpec() instanceof
DistributionSpecReplicated) {
planContext.setBroadcastJoin();
}
+ // The children of an aggregate are optimized before its final cost is
computed
+ // (see CostAndEnforcerJob), so the memo tells whether the translator
fuses it.
+ if (groupExpression.getPlan() instanceof PhysicalHashAggregate &&
childrenProperties.size() == 1
Review Comment:
[P2] Keep the bucketed cost gate valid after shuffle-key pruning. A full-key
`GlobalAgg(GROUP BY a,b) -> Distribute(HASH(a,b)) -> OlapScan` passes this gate
and receives the bucketed discount and one-phase exemption. The default
`ShuffleKeyPruner` can later shorten the distribute to `HASH(a)` when `a` is
balanced. The translator's exact-key check then rejects fusion, leaving a
regular aggregate above an exchange of raw scan rows even though the plan was
costed as fused. Preserve the full keys for a selected bucketed aggregate or
re-evaluate its cost and eligibility after pruning; cover the default plan
choice with analyzed two-key stats.
##########
be/src/exec/operator/bucketed_aggregation_source_operator.cpp:
##########
@@ -0,0 +1,790 @@
+// 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/bucketed_aggregation_sink_operator.h"
+#include "exec/operator/inline_count.h"
+#include "exec/operator/operator.h"
+#include "exprs/vectorized_agg_fn.h"
+#include "runtime/runtime_profile.h"
+#include "runtime/thread_context.h"
+#include "util/debug_points.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 the inline counters and writes the sum back through
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) {
+ inline_count_add(dst_ref, inline_count_get(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");
+ _memory_usage_merge_arena =
+ ADD_COUNTER(custom_profile(), "MemoryUsageMergeArena",
TUnit::BYTES);
+ _memory_usage_merged_hash_tables =
+ ADD_COUNTER(custom_profile(), "MemoryUsageMergedHashTables",
TUnit::BYTES);
+
+ return Status::OK();
+}
+
+void BucketedAggLocalState::_update_memusage(Arena& merge_arena) {
+ int64_t arena_memory_usage = merge_arena.size();
+ COUNTER_SET(_memory_usage_merge_arena, arena_memory_usage);
+ COUNTER_SET(_memory_usage_merged_hash_tables, _hash_table_merge_growth);
+ COUNTER_SET(_memory_used_counter, arena_memory_usage +
_hash_table_merge_growth);
+}
+
+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->notify_state_changed();
+ }
+
+ 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);
+ }
+ }
+}
+
+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;
+ const int64_t dst_bytes_before =
dst_data.get_buffer_size_in_bytes();
+
+ // 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;
Review Comment:
[P2] Avoid waking every source for empty bucket merges. `_merge_bucket`
increments `merged_count` for each finished sink even when this bucket has no
mapped entries or null key. With sparse grouped input and several sink/source
tasks, each intermediate sink completion therefore calls
`notify_state_changed()` up to 256 times, and each call visits every source
dependency; the awakened tasks rescan buckets despite no new aggregate data.
Count only buckets with transferred state and coalesce the wakeup after the
scan while preserving the generation check for real state changes.
--
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]