This is an automated email from the ASF dual-hosted git repository. yiguolei pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/doris.git
commit 37352da41c13a161d7a87f178dcd062ff7624d4f 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]
