This is an automated email from the ASF dual-hosted git repository.

mrhhsg pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 3ecd75c8c34 [fix](be) Use per-source-instance arenas when merging 
bucketed hash agg states (#68564)
3ecd75c8c34 is described below

commit 3ecd75c8c34d3477dde0d647e31cac3feffc98f1
Author: Jerry Hu <[email protected]>
AuthorDate: Tue Sep 29 18:37:16 2026 +0800

    [fix](be) Use per-source-instance arenas when merging bucketed hash agg 
states (#68564)
    
    ### What problem does this PR solve?
    
    Issue Number: None
    
    Related PR: #61495
    
    Problem Summary: The source side of bucketed hash aggregation merges the
    per-sink-instance hash tables bucket by bucket. The per-bucket CAS lock
    (`merge_in_progress`) only serializes work on the same bucket, so two
    source
    tasks can merge two different buckets at the same time. Both of them
    passed
    `*src_inst.arena` (the arena of the sink instance being merged) to
    `merge_agg_states` / `merge_null_key`, i.e. the same non-thread-safe
    `Arena`
    was used concurrently by several threads. Aggregate functions whose
    merge
    allocates from the arena (for example `collect_set` on strings through
    `arena.insert`, or DISTINCT on strings) then got overlapping memory,
    which
    corrupts the merged states: wrong results and potentially BE crashes.
    
    Reproduce on a single BE with `enable_bucketed_hash_agg=true`,
    `parallel_pipeline_task_num=8` and a table whose group keys are present
    in all
    tablets:
    
        SELECT k, collect_set(s) FROM t GROUP BY k;
    
    Before the fix the total number of collected elements was randomly lower
    than
    expected (e.g. 299431 / 299484 instead of 300000) and differed between
    runs.
    After the fix the result is stable and equal to the non-bucketed plan.
    
    Fix: give every source instance its own merge `Arena`, kept in
    `BucketedAggSharedState::source_merge_arenas` and indexed by the source
    task
    index. A source task runs on one thread at a time, so its arena is never
    used
    concurrently. The arenas live as long as the shared state because merged
    states may point into them and any source instance may output a bucket.
    One
    arena per source instance (instead of one per bucket) keeps the retained
    memory proportional to the source parallelism rather than to the 256
    buckets,
    since each arena allocates at least a 4 KiB chunk once used. The final
    null-key merge keeps using the merge target's arena since it runs
    single-threaded after all buckets are done.
    
    ### Release note
    
    Fix wrong results or crashes of bucketed hash aggregation with aggregate
    functions whose merge allocates memory, such as collect_set on strings.
    
    ### Check List (For Author)
    
    - Test:
    - Regression test: added
    query_p0/aggregate/collect_set_bucketed_agg_merge
    (reproduced the wrong result without the fix, passes with it); also ran
          percentile_bucketed_agg_merge and agg_strategy/bucketed_hash_agg
    - Unit Test: added BucketedAggSharedStateTest for the
    per-source-instance
          merge arenas
    - Behavior changed: No
    - Does this need documentation: No
---
 .../bucketed_aggregation_source_operator.cpp       |  8 ++-
 be/src/exec/pipeline/dependency.cpp                |  4 ++
 be/src/exec/pipeline/dependency.h                  |  9 +++
 .../pipeline/bucketed_agg_shared_state_test.cpp    | 44 +++++++++++++
 .../aggregate/collect_set_bucketed_agg_merge.out   | 19 ++++++
 .../collect_set_bucketed_agg_merge.groovy          | 74 ++++++++++++++++++++++
 6 files changed, 156 insertions(+), 2 deletions(-)

diff --git a/be/src/exec/operator/bucketed_aggregation_source_operator.cpp 
b/be/src/exec/operator/bucketed_aggregation_source_operator.cpp
index 993f1868e17..44f35b06836 100644
--- a/be/src/exec/operator/bucketed_aggregation_source_operator.cpp
+++ b/be/src/exec/operator/bucketed_aggregation_source_operator.cpp
@@ -205,6 +205,10 @@ 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];
@@ -291,7 +295,7 @@ int BucketedAggLocalState::_merge_bucket(int bucket, int 
merge_target) {
                                                                             
.aggregate_evaluators,
                                                                     
shared_state
                                                                             
.offsets_of_aggregate_states,
-                                                                    
*src_inst.arena);
+                                                                    
merge_arena);
                                                         }
                                                     });
 
@@ -300,7 +304,7 @@ int BucketedAggLocalState::_merge_bucket(int bucket, int 
merge_target) {
                                                             
shared_state.aggregate_evaluators,
                                                             shared_state
                                                                     
.offsets_of_aggregate_states,
-                                                            *src_inst.arena);
+                                                            merge_arena);
 
                                                     // Mark this instance as 
merged for this bucket.
                                                     
bs.merged_instances[inst_idx] = true;
diff --git a/be/src/exec/pipeline/dependency.cpp 
b/be/src/exec/pipeline/dependency.cpp
index f808e6b10e8..85ca4f0c21b 100644
--- a/be/src/exec/pipeline/dependency.cpp
+++ b/be/src/exec/pipeline/dependency.cpp
@@ -272,6 +272,10 @@ Status BucketedAggSharedState::init_instances(int 
num_instances,
         for (auto& bs : bucket_states) {
             bs.merged_instances.resize(num_instances, false);
         }
+        source_merge_arenas.resize(source_deps.size());
+        for (auto& arena : source_merge_arenas) {
+            arena = std::make_unique<Arena>();
+        }
         _init_status = metadata_init();
     });
     return _init_status;
diff --git a/be/src/exec/pipeline/dependency.h 
b/be/src/exec/pipeline/dependency.h
index 5802a8ddefc..714b9c10a52 100644
--- a/be/src/exec/pipeline/dependency.h
+++ b/be/src/exec/pipeline/dependency.h
@@ -481,6 +481,15 @@ public:
     /// Per-bucket merge state. Indexed by bucket id [0, 256).
     std::array<BucketMergeState, BUCKETED_AGG_NUM_BUCKETS> bucket_states;
 
+    /// Arenas for memory allocated by aggregate function merges on the source 
side.
+    /// One per source instance, indexed by the source task idx and sized in 
init_instances().
+    /// Sink instance arenas cannot be used there because different buckets 
are merged
+    /// concurrently by different source instances and Arena is not 
thread-safe, while a
+    /// source instance runs on one thread at a time. Merged states may point 
into these
+    /// arenas and may be output by any source instance, so they live as long 
as the
+    /// shared state.
+    std::vector<std::unique_ptr<Arena>> source_merge_arenas;
+
     // Aggregate function metadata (shared, read-only after init).
     std::vector<AggFnEvaluator*> aggregate_evaluators;
     VExprContextSPtrs probe_expr_ctxs;
diff --git a/be/test/exec/pipeline/bucketed_agg_shared_state_test.cpp 
b/be/test/exec/pipeline/bucketed_agg_shared_state_test.cpp
new file mode 100644
index 00000000000..39ecaba7a67
--- /dev/null
+++ b/be/test/exec/pipeline/bucketed_agg_shared_state_test.cpp
@@ -0,0 +1,44 @@
+// 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 <gtest/gtest.h>
+
+#include "exec/common/agg_utils.h"
+#include "exec/pipeline/dependency.h"
+
+namespace doris {
+
+// Source instances merge different buckets concurrently, so each of them 
needs its own merge
+// arena. The arenas are owned by the shared state and are independent of the 
sink count and
+// of the number of buckets.
+TEST(BucketedAggSharedStateTest, InitInstancesCreatesOneMergeArenaPerSource) {
+    BucketedAggSharedState state;
+    state.create_source_dependencies(3, 0, 0, "BUCKETED_AGG_SOURCE");
+    ASSERT_TRUE(state.init_instances(2, [] { return Status::OK(); }).ok());
+
+    ASSERT_EQ(state.source_merge_arenas.size(), 3);
+    for (size_t i = 0; i < state.source_merge_arenas.size(); ++i) {
+        ASSERT_NE(state.source_merge_arenas[i], nullptr);
+        // No memory is reserved until a merge allocates from the arena.
+        EXPECT_EQ(state.source_merge_arenas[i]->size(), 0);
+        for (size_t j = 0; j < i; ++j) {
+            EXPECT_NE(state.source_merge_arenas[i].get(), 
state.source_merge_arenas[j].get());
+        }
+    }
+}
+
+} // namespace doris
diff --git 
a/regression-test/data/query_p0/aggregate/collect_set_bucketed_agg_merge.out 
b/regression-test/data/query_p0/aggregate/collect_set_bucketed_agg_merge.out
new file mode 100644
index 00000000000..6781b8802ef
--- /dev/null
+++ b/regression-test/data/query_p0/aggregate/collect_set_bucketed_agg_merge.out
@@ -0,0 +1,19 @@
+-- This file is automatically generated. You should know what you did if you 
want to edit this
+-- !bucketed --
+3000   300000  300000  -253297445209059746020
+
+-- !bucketed --
+3000   300000  300000  -253297445209059746020
+
+-- !bucketed --
+3000   300000  300000  -253297445209059746020
+
+-- !bucketed --
+3000   300000  300000  -253297445209059746020
+
+-- !bucketed --
+3000   300000  300000  -253297445209059746020
+
+-- !no_bucketed --
+3000   300000  300000  -253297445209059746020
+
diff --git 
a/regression-test/suites/query_p0/aggregate/collect_set_bucketed_agg_merge.groovy
 
b/regression-test/suites/query_p0/aggregate/collect_set_bucketed_agg_merge.groovy
new file mode 100644
index 00000000000..efd850ebcf8
--- /dev/null
+++ 
b/regression-test/suites/query_p0/aggregate/collect_set_bucketed_agg_merge.groovy
@@ -0,0 +1,74 @@
+// 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.
+
+// Bucketed hash aggregation merges the per-instance states of different 
buckets
+// concurrently in several source instances. Merges that allocate memory 
(collect_set
+// on strings, DISTINCT on strings) must not share an arena across those 
instances.
+suite("collect_set_bucketed_agg_merge") {
+    sql "set enable_bucketed_hash_agg=true"
+    sql "set be_number_for_test=1"
+    sql "set agg_phase=1"
+    sql "set parallel_pipeline_task_num=8"
+    sql "set bucketed_agg_min_input_rows=0"
+    sql "set bucketed_agg_max_group_keys=0"
+    sql "set bucketed_agg_high_card_threshold=1.0"
+
+    sql "DROP TABLE IF EXISTS collect_set_bucketed_agg_merge_t"
+    sql """
+        CREATE TABLE collect_set_bucketed_agg_merge_t (
+            id INT NOT NULL,
+            k INT NOT NULL,
+            s VARCHAR(128) NOT NULL
+        )
+        DUPLICATE KEY(id)
+        DISTRIBUTED BY HASH(id) BUCKETS 16
+        PROPERTIES ('replication_num' = '1')
+    """
+    // every group key has rows in every tablet, so every sink instance holds
+    // states for keys spread over all buckets
+    sql """
+        INSERT INTO collect_set_bucketed_agg_merge_t
+        SELECT number, number % 3000,
+               concat('value_', number % 13, '_', repeat('x', number % 37))
+        FROM numbers("number" = "300000")
+    """
+
+    String query = """
+        SELECT k, collect_set(s) cs, count(DISTINCT s) cnt
+        FROM collect_set_bucketed_agg_merge_t GROUP BY k
+    """
+    explain {
+        sql query
+        contains("BUCKETED AGGREGATE")
+    }
+
+    for (int i = 0; i < 5; i++) {
+        order_qt_bucketed """
+            SELECT count(*), sum(size(cs)), sum(cnt),
+                   sum(cast(murmur_hash3_64(concat(k, ':', 
array_join(array_sort(cs), ','))) AS LARGEINT))
+            FROM (${query}) q
+        """
+    }
+
+    // control: the same query without bucketed hash aggregation
+    sql "set enable_bucketed_hash_agg=false"
+    order_qt_no_bucketed """
+        SELECT count(*), sum(size(cs)), sum(cnt),
+               sum(cast(murmur_hash3_64(concat(k, ':', 
array_join(array_sort(cs), ','))) AS LARGEINT))
+        FROM (${query}) q
+    """
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to