This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 9f181d9ee0d [fix](be) Backport unlimited topn expansion rates to
branch-4.1 (#68545)
9f181d9ee0d is described below
commit 9f181d9ee0d89f709026b4c9d0bb671b913927c3
Author: HappenLee <[email protected]>
AuthorDate: Tue Sep 29 17:40:40 2026 +0800
[fix](be) Backport unlimited topn expansion rates to branch-4.1 (#68545)
### What problem does this PR solve?
Issue Number: N/A
Related PR: #68246
Backport #68246 to `branch-4.1`. A zero `space_expand_rate` previously
produced a zero candidate capacity, so partial-state serialization
discarded every TOPN candidate and two-stage aggregation could return an
empty result. Treat non-positive rates as unlimited intermediate
candidate retention for `topn`, `topn_array`, and `topn_weighted`;
`top_num` still limits the final result. Positive and default rates
retain their existing behavior.
Branch-specific conflict resolution:
- Keep the production fix, standalone TOPN unit tests, regression suite,
and the expected output generated in the original PR.
- Omit modifications to `agg_state_parameters_test.cpp` and
`test_agg_state_parameters.groovy`: neither file nor the associated
aggregate-state parameter validation exists in branch-4.1. This backport
does not introduce that separate master feature.
- Remove `SET enable_bucketed_hash_agg = false` from the regression
suite because branch-4.1 has no such session variable. Keep coverage for
aggregation phases 1 and 2.
### Release note
`topn`, `topn_array`, and `topn_weighted` interpret a non-positive
`space_expand_rate` as unlimited intermediate candidate retention. The
final result remains limited by `top_num`. Retaining all distinct
candidates can increase intermediate-state memory and network traffic.
### Check List (For Author)
- Test: Unit Test / Regression test included; local execution not
completed
- Repository clang-format 16 formatting and check scripts passed.
- Diff whitespace check passed with the original generated `.out`
trailing blank line preserved.
- Attempted `./run-be-ut.sh -j 48 --run
--filter='AggregateFunctionTopN*.*:*/AggregateFunctionTopN*.*:AggTest.topn*'`.
Stopped during dependency preparation; the isolated checkout has no
installed third-party dependency bundle. No local 4.1 unit-test pass is
claimed.
- Regression tests were not run locally against a 4.1 cluster. Expected
output is copied unchanged from #68246.
- Branch-4.1 does not contain `build-support/check-build-hygiene.sh`; no
successful hygiene or clang-tidy run is claimed.
- Behavior changed: Yes. Non-positive expansion rates retain all
intermediate candidates.
- Does this need documentation: Yes. The function documentation should
describe non-positive expansion rates; the documentation follow-up noted
by #68246 remains applicable.
---
be/src/exprs/aggregate/aggregate_function_topn.h | 3 +-
be/test/exprs/aggregate/agg_topn_test.cpp | 125 +++++++++++++++++++++
.../agg_function/topn/topn_unlimited.out | 91 +++++++++++++++
.../agg_function/topn/topn_unlimited.groovy | 75 +++++++++++++
4 files changed, 293 insertions(+), 1 deletion(-)
diff --git a/be/src/exprs/aggregate/aggregate_function_topn.h
b/be/src/exprs/aggregate/aggregate_function_topn.h
index 77c9e260d4a..aae1bc182c7 100644
--- a/be/src/exprs/aggregate/aggregate_function_topn.h
+++ b/be/src/exprs/aggregate/aggregate_function_topn.h
@@ -60,7 +60,8 @@ struct AggregateFunctionTopNData {
using DataType = typename PrimitiveTypeTraits<T>::CppType;
void set_paramenters(int input_top_num, int space_expand_rate = 50) {
top_num = input_top_num;
- capacity = (uint64_t)top_num * space_expand_rate;
+ // Non-positive expansion rates retain all candidates during
serialization and merging.
+ capacity = space_expand_rate <= 0 ? UINT64_MAX : (uint64_t)top_num *
space_expand_rate;
}
void add(const StringRef& value, const UInt64& increment = 1) {
diff --git a/be/test/exprs/aggregate/agg_topn_test.cpp
b/be/test/exprs/aggregate/agg_topn_test.cpp
new file mode 100644
index 00000000000..324bb7272d9
--- /dev/null
+++ b/be/test/exprs/aggregate/agg_topn_test.cpp
@@ -0,0 +1,125 @@
+// 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 <cstdint>
+#include <string>
+
+#include "core/column/column_string.h"
+#include "core/column/column_vector.h"
+#include "core/string_buffer.hpp"
+#include "exprs/aggregate/aggregate_function_topn.h"
+
+namespace doris {
+namespace {
+
+template <PrimitiveType T>
+AggregateFunctionTopNData<T> round_trip(const AggregateFunctionTopNData<T>&
state) {
+ auto column = ColumnString::create();
+ BufferWritable writer(*column);
+ state.write(writer);
+ writer.commit();
+ BufferReadable reader(column->get_data_at(0));
+ AggregateFunctionTopNData<T> result;
+ result.read(reader);
+ return result;
+}
+
+class AggregateFunctionTopNUnlimitedTest : public testing::TestWithParam<int>
{};
+
+TEST_P(AggregateFunctionTopNUnlimitedTest, SerializeAndMergeStrings) {
+ AggregateFunctionTopNData<TYPE_STRING> lhs;
+ AggregateFunctionTopNData<TYPE_STRING> rhs;
+ lhs.set_paramenters(1, GetParam());
+ rhs.set_paramenters(1, GetParam());
+ // The global winner is not the most frequent value in either partial
state.
+ lhs.add(std::string("a"), 3);
+ lhs.add(std::string("winner"), 2);
+ rhs.add(std::string("b"), 3);
+ rhs.add(std::string("winner"), 2);
+
+ auto partial = round_trip(lhs);
+ EXPECT_EQ(partial.counter_map, lhs.counter_map);
+ AggregateFunctionTopNData<TYPE_STRING> merged;
+ merged.merge(partial);
+ merged = round_trip(merged);
+ merged.merge(round_trip(rhs));
+ merged = round_trip(merged);
+
+ ASSERT_EQ(merged.counter_map.size(), 3);
+ EXPECT_EQ(merged.counter_map.at("a"), 3);
+ EXPECT_EQ(merged.counter_map.at("b"), 3);
+ EXPECT_EQ(merged.counter_map.at("winner"), 4);
+ EXPECT_EQ(merged.get(), R"({"winner":4})");
+
+ AggregateFunctionTopNData<TYPE_STRING> empty;
+ merged.merge(round_trip(empty));
+ EXPECT_EQ(merged.get(), R"({"winner":4})");
+
+ merged.reset();
+ EXPECT_EQ(round_trip(merged).get(), "{}");
+ merged.set_paramenters(1, GetParam());
+ merged.add(std::string("new"), 7);
+ EXPECT_EQ(round_trip(merged).get(), R"({"new":7})");
+}
+
+TEST_P(AggregateFunctionTopNUnlimitedTest, SerializeAndMergeWeightedIntegers) {
+ AggregateFunctionTopNData<TYPE_INT> lhs;
+ AggregateFunctionTopNData<TYPE_INT> rhs;
+ lhs.set_paramenters(2, GetParam());
+ rhs.set_paramenters(2, GetParam());
+ lhs.add(1, 10);
+ lhs.add(2, 7);
+ lhs.add(3, 6);
+ rhs.add(4, 11);
+ rhs.add(5, 8);
+ rhs.add(3, 6);
+
+ AggregateFunctionTopNData<TYPE_INT> merged;
+ merged.merge(round_trip(lhs));
+ merged.merge(round_trip(rhs));
+ merged = round_trip(merged);
+ ASSERT_EQ(merged.counter_map.size(), 5);
+ EXPECT_EQ(merged.counter_map.at(3), 12);
+ auto result = ColumnInt32::create();
+ merged.insert_result_into(*result);
+ ASSERT_EQ(result->size(), 2);
+ EXPECT_EQ(result->get_element(0), 3);
+ EXPECT_EQ(result->get_element(1), 4);
+}
+
+INSTANTIATE_TEST_SUITE_P(NonPositiveRates, AggregateFunctionTopNUnlimitedTest,
+ testing::Values(0, -1, INT32_MIN));
+
+TEST(AggregateFunctionTopNTest, PositiveRateStillLimitsSerializedCandidates) {
+ AggregateFunctionTopNData<TYPE_INT> state;
+ state.set_paramenters(1, 2);
+ state.add(1, 3);
+ state.add(2, 2);
+ state.add(3, 1);
+ auto partial = round_trip(state);
+ ASSERT_EQ(partial.counter_map.size(), 2);
+ EXPECT_EQ(partial.counter_map.at(1), 3);
+ EXPECT_EQ(partial.counter_map.at(2), 2);
+
+ state.set_paramenters(1);
+ EXPECT_EQ(round_trip(state).counter_map, state.counter_map);
+}
+
+} // namespace
+} // namespace doris
diff --git
a/regression-test/data/nereids_function_p0/agg_function/topn/topn_unlimited.out
b/regression-test/data/nereids_function_p0/agg_function/topn/topn_unlimited.out
new file mode 100644
index 00000000000..888ec8e5112
--- /dev/null
+++
b/regression-test/data/nereids_function_p0/agg_function/topn/topn_unlimited.out
@@ -0,0 +1,91 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !unlimited --
+{"a":3,"x":2} ["a", "x"] [1, 4] ["b", "y"] [2, 5]
+
+-- !grouped --
+1 {"a":3,"b":2} ["a", "b"] [1, 2] ["b", "a"] [2, 1]
+2 {"x":2,"y":1} ["x", "y"] [4, 5] ["y", "x"] [5, 4]
+3 \N \N \N \N \N
+
+-- !empty --
+\N \N \N
+
+-- !constant_input --
+{"a":2} ["a"]
+
+-- !unlimited --
+{"a":3,"x":2} ["a", "x"] [1, 4] ["b", "y"] [2, 5]
+
+-- !grouped --
+1 {"a":3,"b":2} ["a", "b"] [1, 2] ["b", "a"] [2, 1]
+2 {"x":2,"y":1} ["x", "y"] [4, 5] ["y", "x"] [5, 4]
+3 \N \N \N \N \N
+
+-- !empty --
+\N \N \N
+
+-- !constant_input --
+{"a":2} ["a"]
+
+-- !unlimited --
+{"a":3,"x":2} ["a", "x"] [1, 4] ["b", "y"] [2, 5]
+
+-- !grouped --
+1 {"a":3,"b":2} ["a", "b"] [1, 2] ["b", "a"] [2, 1]
+2 {"x":2,"y":1} ["x", "y"] [4, 5] ["y", "x"] [5, 4]
+3 \N \N \N \N \N
+
+-- !empty --
+\N \N \N
+
+-- !constant_input --
+{"a":2} ["a"]
+
+-- !positive_and_default --
+{"a":3,"x":2} {"a":3,"x":2} [1, 4] [1, 4] [2, 5] [2, 5]
+
+-- !unlimited --
+{"a":3,"x":2} ["a", "x"] [1, 4] ["b", "y"] [2, 5]
+
+-- !grouped --
+1 {"a":3,"b":2} ["a", "b"] [1, 2] ["b", "a"] [2, 1]
+2 {"x":2,"y":1} ["x", "y"] [4, 5] ["y", "x"] [5, 4]
+3 \N \N \N \N \N
+
+-- !empty --
+\N \N \N
+
+-- !constant_input --
+{"a":2} ["a"]
+
+-- !unlimited --
+{"a":3,"x":2} ["a", "x"] [1, 4] ["b", "y"] [2, 5]
+
+-- !grouped --
+1 {"a":3,"b":2} ["a", "b"] [1, 2] ["b", "a"] [2, 1]
+2 {"x":2,"y":1} ["x", "y"] [4, 5] ["y", "x"] [5, 4]
+3 \N \N \N \N \N
+
+-- !empty --
+\N \N \N
+
+-- !constant_input --
+{"a":2} ["a"]
+
+-- !unlimited --
+{"a":3,"x":2} ["a", "x"] [1, 4] ["b", "y"] [2, 5]
+
+-- !grouped --
+1 {"a":3,"b":2} ["a", "b"] [1, 2] ["b", "a"] [2, 1]
+2 {"x":2,"y":1} ["x", "y"] [4, 5] ["y", "x"] [5, 4]
+3 \N \N \N \N \N
+
+-- !empty --
+\N \N \N
+
+-- !constant_input --
+{"a":2} ["a"]
+
+-- !positive_and_default --
+{"a":3,"x":2} {"a":3,"x":2} [1, 4] [1, 4] [2, 5] [2, 5]
+
diff --git
a/regression-test/suites/nereids_function_p0/agg_function/topn/topn_unlimited.groovy
b/regression-test/suites/nereids_function_p0/agg_function/topn/topn_unlimited.groovy
new file mode 100644
index 00000000000..89546bfb5fa
--- /dev/null
+++
b/regression-test/suites/nereids_function_p0/agg_function/topn/topn_unlimited.groovy
@@ -0,0 +1,75 @@
+// 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("topn_unlimited") {
+ sql "DROP TABLE IF EXISTS test_topn_unlimited"
+ sql """
+ CREATE TABLE test_topn_unlimited (
+ g INT,
+ s STRING,
+ v INT,
+ w BIGINT
+ ) DISTRIBUTED BY HASH(v) BUCKETS 3
+ PROPERTIES ("replication_num" = "1")
+ """
+ sql """
+ INSERT INTO test_topn_unlimited VALUES
+ (1, 'a', 1, 1), (1, 'b', 2, 10), (1, 'a', 1, 2),
+ (1, 'c', 3, 3), (1, 'a', 1, 1), (1, 'b', 2, 2),
+ (2, 'x', 4, 1), (2, 'y', 5, 5), (2, 'x', 4, 2),
+ (3, NULL, NULL, 1)
+ """
+
+ sql "SET parallel_pipeline_task_num = 1"
+ for (def phase : [1, 2]) {
+ sql "SET agg_phase = ${phase}"
+ // Non-positive rates retain every candidate, including after partial
serialization.
+ for (def rate : [0, -1, -2147483648]) {
+ order_qt_unlimited """
+ SELECT topn(s, 2, ${rate}),
+ topn_array(s, 2, ${rate}),
+ topn_array(v, 2, ${rate}),
+ topn_weighted(s, w, 2, ${rate}),
+ topn_weighted(v, w, 2, ${rate})
+ FROM test_topn_unlimited
+ """
+ order_qt_grouped """
+ SELECT g, topn(s, 2, ${rate}),
+ topn_array(s, 2, ${rate}),
+ topn_array(v, 2, ${rate}),
+ topn_weighted(s, w, 2, ${rate}),
+ topn_weighted(v, w, 2, ${rate})
+ FROM test_topn_unlimited GROUP BY g
+ """
+ order_qt_empty """
+ SELECT topn(s, 1, ${rate}), topn_array(v, 1, ${rate}),
+ topn_weighted(v, w, 1, ${rate})
+ FROM test_topn_unlimited WHERE g = 4
+ """
+ order_qt_constant_input """
+ SELECT topn(s, 1, ${rate}), topn_array(s, 1, ${rate})
+ FROM (SELECT 'a' AS s UNION ALL SELECT 'b' UNION ALL SELECT
'a') t
+ """
+ }
+ order_qt_positive_and_default """
+ SELECT topn(s, 2), topn(s, 2, 50),
+ topn_array(v, 2), topn_array(v, 2, 50),
+ topn_weighted(v, w, 2), topn_weighted(v, w, 2, 50)
+ FROM test_topn_unlimited
+ """
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]