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

HappenLee 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 1cc0456c028 [fix](agg_state) Enable pre-aggregation for aggregate 
state merge and union (#68496)
1cc0456c028 is described below

commit 1cc0456c02839da44c64c08bba24c691965e2d43
Author: HappenLee <[email protected]>
AuthorDate: Thu Oct 8 17:18:22 2026 +0800

    [fix](agg_state) Enable pre-aggregation for aggregate state merge and union 
(#68496)
    
    ### What problem does this PR solve?
    
    Issue Number: N/A
    
    Queries such as `SELECT k, max_by_merge(s) FROM t GROUP BY k` over a
    `GENERIC` aggregate-state column report `PREAGGREGATION: OFF` because
    the value-column checker does not recognize aggregate state combinators.
    The scan therefore retains the storage merge requirement even when the
    query aggregate can merge the partial states.
    
    Handle both `MergeCombinator` and `UnionCombinator` through a shared
    check for `GENERIC` columns and matching nested function names. Keep the
    existing bare-slot, filter, grouping, join and sibling-aggregate
    restrictions. Cover multiple rowsets, `max_by` and non-idempotent `sum`
    states, NULL/empty inputs, aliases, grouping and aggregate phases, plus
    negative plan cases.
    
    Pre-aggregation can change partial-state merge order. Normalize
    unordered array outputs with `array_sort` and concatenated string
    outputs with split/sort/join in six existing regression suites; the
    fixtures do not contain delimiter characters. Keep MV plan assertions on
    the normalized queries and regenerate their expected outputs. Keep
    `test_agg_state_map` active with unique map keys within each group,
    retaining NULL key/value coverage. Sort keys and reorder values by those
    same keys so each association remains intact. Assert that the normalized
    query enables pre-aggregation and compare its generated expected result.
    
    Document the existing MAP_AGG contract: retain an already-seen key's
    value and ignore subsequent values. Encounter order is execution order
    and does not guarantee the earliest inserted record or a repeatable
    winner. Documentation: https://github.com/apache/doris-website/pull/4175
    
    ### Release note
    
    Enable scan pre-aggregation for compatible aggregate state merge and
    union functions over GENERIC columns.
    
    ### Check List (For Author)
    
    - Test:
    - [x] Follow-up: expected outputs for all seven affected suites were
    generated with `-forceGenOut`; the restored MAP_AGG_STATE suite, the six
    normalized suites and `agg_state_preagg` passed a separate normal
    comparison run (8/8).
    - [x] The changed suites cover map/array/group-concat states, MV
    rewrites, outfile queries and the cloud-directory serialization case.
    The local run uses a shared-nothing cluster; cloud deployment validation
    is left to CI.
    - [x] FE build and Checkstyle passed; regression source/output
    whitespace checks passed.
    - [x] Previous implementation validation: `SetPreAggStatusTest` 13/13;
    `agg_state_preagg`, `set_preagg`, `test_agg_state` 3/3. This follow-up
    changes only regression files.
    - Local follow-up validation used the newly built FE and existing ASAN
    BE `dd3930c6844` at `be_exec_version=14`. It does not replace CI
    coverage of the current BE/default execution version. No BE code changes
    or performance measurements are included.
    - Behavior changed:
    - [x] Yes. Compatible merge/union calls permit scan pre-aggregation.
    Existing incompatible paths remain OFF. Unordered aggregates can expose
    a different element order or duplicate-key winner.
    - Does this need documentation?
    - [x] Yes. Clarify MAP_AGG duplicate-key selection in the linked
    documentation PR.
    
    ### Check List (For Reviewer who merge this PR)
    
    - [ ] Confirm the release note
    - [ ] Confirm test cases
    - [ ] Confirm document
    - [ ] Add branch pick label
---
 .../nereids/rules/rewrite/SetPreAggStatus.java     |  29 +++++
 .../nereids/rules/rewrite/SetPreAggStatusTest.java |  51 ++++++++
 .../different_serialize/different_serialize.out    |   6 +-
 .../data/datatype_p0/agg_state/array/array.out     |   2 +-
 .../group_concat/test_agg_state_group_concat.out   |  10 +-
 .../data/datatype_p0/agg_state/map/map.out         |   4 +-
 .../diffrent_serialize/diffrent_serialize.out      |   2 +-
 .../set_preagg/agg_state_preagg.out                |  65 ++++++++++
 .../outfile/agg_state/test_outfile_agg_state.out   |   4 +-
 .../agg_state_array/test_outfile_agg_array.out     |   2 +-
 .../different_serialize/different_serialize.groovy |  12 +-
 .../datatype_p0/agg_state/array/array.groovy       |   2 +-
 .../test_agg_state_group_concat.groovy             |  10 +-
 .../suites/datatype_p0/agg_state/map/map.groovy    |  19 ++-
 .../diffrent_serialize/diffrent_serialize.groovy   |  16 +--
 .../set_preagg/agg_state_preagg.groovy             | 137 +++++++++++++++++++++
 .../agg_state/test_outfile_agg_state.groovy        |   4 +-
 .../agg_state_array/test_outfile_agg_array.groovy  |   4 +-
 18 files changed, 336 insertions(+), 43 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SetPreAggStatus.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SetPreAggStatus.java
index f81b6159200..711b7b3fa08 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SetPreAggStatus.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/SetPreAggStatus.java
@@ -40,6 +40,9 @@ import 
org.apache.doris.nereids.trees.expressions.functions.agg.HllUnionAgg;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Max;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Min;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Sum;
+import 
org.apache.doris.nereids.trees.expressions.functions.combinator.Combinator;
+import 
org.apache.doris.nereids.trees.expressions.functions.combinator.MergeCombinator;
+import 
org.apache.doris.nereids.trees.expressions.functions.combinator.UnionCombinator;
 import 
org.apache.doris.nereids.trees.expressions.functions.scalar.GroupingScalarFunction;
 import org.apache.doris.nereids.trees.expressions.functions.scalar.If;
 import org.apache.doris.nereids.trees.expressions.literal.NullLiteral;
@@ -55,6 +58,7 @@ import 
org.apache.doris.nereids.trees.plans.logical.LogicalProject;
 import org.apache.doris.nereids.trees.plans.logical.LogicalRepeat;
 import org.apache.doris.nereids.trees.plans.visitor.CustomRewriter;
 import org.apache.doris.nereids.trees.plans.visitor.DefaultPlanRewriter;
+import org.apache.doris.nereids.types.AggStateType;
 import org.apache.doris.nereids.util.ExpressionUtils;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.qe.SessionVariable;
@@ -705,6 +709,31 @@ public class SetPreAggStatus extends 
DefaultPlanRewriter<Stack<SetPreAggStatus.P
                         .off(String.format("%s is not supported.", 
aggregateFunction.toSql()));
             }
 
+            @Override
+            public PreAggStatus visitMergeCombinator(MergeCombinator 
combinator, AggregateType aggregateType) {
+                return checkAggStateCombinator(combinator, aggregateType);
+            }
+
+            @Override
+            public PreAggStatus visitUnionCombinator(UnionCombinator 
combinator, AggregateType aggregateType) {
+                return checkAggStateCombinator(combinator, aggregateType);
+            }
+
+            private PreAggStatus checkAggStateCombinator(AggregateFunction 
aggregateFunction,
+                    AggregateType aggregateType) {
+                // GENERIC merges stored states with the same aggregate 
function. A matching
+                // merge/union can consume the partial states directly; 
REPLACE cannot.
+                // The caller requires a bare value slot, and the combinator 
builder derives
+                // the nested argument types and nullability from that slot's 
AggStateType.
+                AggStateType stateType = (AggStateType) 
aggregateFunction.child(0).getDataType();
+                String functionName = ((Combinator) 
aggregateFunction).getNestedFunction().getName();
+                if (aggregateType == AggregateType.GENERIC && 
stateType.getFunctionName().equals(functionName)) {
+                    return PreAggStatus.on();
+                }
+                return PreAggStatus.off(String.format("%s is not match agg 
mode %s or state function %s",
+                        aggregateFunction.toSql(), aggregateType, 
stateType.getFunctionName()));
+            }
+
             @Override
             public PreAggStatus visitMax(Max max, AggregateType aggregateType) 
{
                 if (aggregateType == AggregateType.MAX && !max.isDistinct()) {
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/SetPreAggStatusTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/SetPreAggStatusTest.java
index aa93a831d07..737ce5fab17 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/SetPreAggStatusTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/SetPreAggStatusTest.java
@@ -29,14 +29,19 @@ import 
org.apache.doris.nereids.trees.expressions.SlotReference;
 import 
org.apache.doris.nereids.trees.expressions.functions.agg.AggregateFunction;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Count;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Max;
+import org.apache.doris.nereids.trees.expressions.functions.agg.MaxBy;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Min;
 import org.apache.doris.nereids.trees.expressions.functions.agg.Sum;
+import 
org.apache.doris.nereids.trees.expressions.functions.combinator.MergeCombinator;
+import 
org.apache.doris.nereids.trees.expressions.functions.combinator.StateCombinator;
+import 
org.apache.doris.nereids.trees.expressions.functions.combinator.UnionCombinator;
 import org.apache.doris.nereids.trees.expressions.functions.scalar.If;
 import org.apache.doris.nereids.trees.expressions.functions.scalar.Random;
 import org.apache.doris.nereids.trees.expressions.literal.DoubleLiteral;
 import org.apache.doris.nereids.trees.expressions.literal.IntegerLiteral;
 import org.apache.doris.nereids.trees.plans.PreAggStatus;
 import org.apache.doris.nereids.trees.plans.logical.LogicalOlapScan;
+import org.apache.doris.nereids.types.AggStateType;
 import org.apache.doris.nereids.types.IntegerType;
 
 import com.google.common.collect.ImmutableList;
@@ -104,6 +109,16 @@ class SetPreAggStatusTest {
                 ImmutableList.of("t"), null, column, null, null);
     }
 
+    private static SlotReference aggStateSlot(AggregateFunction nested, 
AggregateType aggregateType) {
+        AggStateType type = (AggStateType) new 
StateCombinator(nested.children(), nested).getDataType();
+        Column column = new Column("s", type.toCatalogDataType(), false, 
aggregateType, null, "");
+        // Column construction normalizes AGG_STATE to GENERIC. Set the mode 
explicitly
+        // to exercise the checker's rejection of non-merging column metadata 
as well.
+        column.setAggregationType(aggregateType, false);
+        return new SlotReference(new ExprId(exprIdCounter++), "s", type, false,
+                ImmutableList.of("t"), null, column, null, null);
+    }
+
     private static PreAggStatus checkAggregateFunctions(
             Set<AggregateFunction> aggregateFuncs, Set<Slot> 
groupingExprsInputSlots, Set<Slot> outputSlots) {
         try {
@@ -139,6 +154,42 @@ class SetPreAggStatusTest {
         return new If(greaterThanZero(key), thenExpr, elseExpr);
     }
 
+    @Test
+    void testAggStateCombinators() {
+        SlotReference value = new SlotReference("value", IntegerType.INSTANCE, 
false);
+        SlotReference order = new SlotReference("order", IntegerType.INSTANCE, 
true);
+        for (AggregateFunction nested : ImmutableList.of(new MaxBy(value, 
order), new Sum(order))) {
+            for (AggregateType aggregateType : 
ImmutableList.of(AggregateType.GENERIC, AggregateType.REPLACE,
+                    AggregateType.REPLACE_IF_NOT_NULL, AggregateType.NONE)) {
+                SlotReference state = aggStateSlot(nested, aggregateType);
+                for (AggregateFunction combinator : ImmutableList.of(
+                        new MergeCombinator(ImmutableList.of(state), nested),
+                        new UnionCombinator(ImmutableList.of(state), nested))) 
{
+                    PreAggStatus status = 
checkAggregateFunctions(Sets.newHashSet(combinator),
+                            Collections.emptySet(), Sets.newHashSet(state));
+                    Assertions.assertEquals(aggregateType == 
AggregateType.GENERIC, status.isOn(),
+                            combinator.toSql() + " over " + aggregateType);
+                }
+            }
+        }
+    }
+
+    @Test
+    void testAggStateCombinatorsRejectDifferentFunctionAndExpression() {
+        SlotReference value = new SlotReference("value", IntegerType.INSTANCE, 
true);
+        AggregateFunction sum = new Sum(value);
+        SlotReference state = aggStateSlot(sum, AggregateType.GENERIC);
+        Expression conditionalState = new If(greaterThanZero(keySlot("k")), 
state, state);
+        for (AggregateFunction combinator : ImmutableList.of(
+                new MergeCombinator(ImmutableList.of(state), new Max(value)),
+                new UnionCombinator(ImmutableList.of(state), new Max(value)),
+                new MergeCombinator(ImmutableList.of(conditionalState), sum),
+                new UnionCombinator(ImmutableList.of(conditionalState), sum))) 
{
+            
Assertions.assertTrue(checkAggregateFunctions(Sets.newHashSet(combinator),
+                    Collections.emptySet(), Sets.newHashSet(state)).isOff(), 
combinator.toSql());
+        }
+    }
+
     @Test
     void testNoAggregateReturnsOff() {
         SlotReference k = keySlot("k");
diff --git 
a/regression-test/data/cloud_p0/different_serialize/different_serialize.out 
b/regression-test/data/cloud_p0/different_serialize/different_serialize.out
index 6bbe497f349..8a370718cad 100644
--- a/regression-test/data/cloud_p0/different_serialize/different_serialize.out
+++ b/regression-test/data/cloud_p0/different_serialize/different_serialize.out
@@ -37,21 +37,21 @@
 -- !select_mv --
 \N     [4]
 -4     [4]
-1      [2, 1, 1]
+1      [1, 1, 2]
 2      [2]
 3      [3]
 
 -- !select_mv --
 \N     [4]
 -4     [4]
-1      [2, 1, 1]
+1      [1, 1, 2]
 2      [2]
 3      [3]
 
 -- !select_mv --
 \N     [4]
 -4     [4]
-1      [2, 1]
+1      [1, 2]
 2      [2]
 3      [3]
 
diff --git a/regression-test/data/datatype_p0/agg_state/array/array.out 
b/regression-test/data/datatype_p0/agg_state/array/array.out
index d1b5a25b4da..dc22cb54e47 100644
--- a/regression-test/data/datatype_p0/agg_state/array/array.out
+++ b/regression-test/data/datatype_p0/agg_state/array/array.out
@@ -1,5 +1,5 @@
 -- This file is automatically generated. You should know what you did if you 
want to edit this
 -- !test --
-1      [2, 1]
+1      [1, 2]
 2      [3]
 
diff --git 
a/regression-test/data/datatype_p0/agg_state/group_concat/test_agg_state_group_concat.out
 
b/regression-test/data/datatype_p0/agg_state/group_concat/test_agg_state_group_concat.out
index 34376e46bb8..4f4ae2f7615 100644
--- 
a/regression-test/data/datatype_p0/agg_state/group_concat/test_agg_state_group_concat.out
+++ 
b/regression-test/data/datatype_p0/agg_state/group_concat/test_agg_state_group_concat.out
@@ -3,20 +3,20 @@
 1      15
 
 -- !group1 --
-1      ccc,bb,a
+1      a,bb,ccc
 
 -- !merge1 --
-ccc,bb,a
+a,bb,ccc
 
 -- !length2 --
 1      15
 
 -- !group2 --
-1      ccc,bb,a
+1      a,bb,ccc
 
 -- !merge2 --
-ccc,bb,a
+a,bb,ccc
 
 -- !union --
-ccc,bb,a
+a,bb,ccc
 
diff --git a/regression-test/data/datatype_p0/agg_state/map/map.out 
b/regression-test/data/datatype_p0/agg_state/map/map.out
index bf64ebfcffa..6a1f4cbee53 100644
--- a/regression-test/data/datatype_p0/agg_state/map/map.out
+++ b/regression-test/data/datatype_p0/agg_state/map/map.out
@@ -1,5 +1,5 @@
 -- This file is automatically generated. You should know what you did if you 
want to edit this
 -- !test --
-1      {null:100, 2:22, 1:11}
-2      {null:400, 4:null, 3:3}
+1      [null, 1, 2, 5, 6]      [100, 1, 2, 11, 22]
+2      [null, 3, 4, 7] [null, 3, null, 400]
 
diff --git 
a/regression-test/data/mv_p0/agg_state/diffrent_serialize/diffrent_serialize.out
 
b/regression-test/data/mv_p0/agg_state/diffrent_serialize/diffrent_serialize.out
index 379496dbaee..7714cf427ba 100644
--- 
a/regression-test/data/mv_p0/agg_state/diffrent_serialize/diffrent_serialize.out
+++ 
b/regression-test/data/mv_p0/agg_state/diffrent_serialize/diffrent_serialize.out
@@ -44,7 +44,7 @@
 3      [3]
 
 -- !select_mv --
-\N     [4, 3]
+\N     [3, 4]
 -4     [4]
 1      [1]
 2      [2]
diff --git 
a/regression-test/data/nereids_rules_p0/set_preagg/agg_state_preagg.out 
b/regression-test/data/nereids_rules_p0/set_preagg/agg_state_preagg.out
new file mode 100644
index 00000000000..c5c2ee7578f
--- /dev/null
+++ b/regression-test/data/nereids_rules_p0/set_preagg/agg_state_preagg.out
@@ -0,0 +1,65 @@
+-- This file is automatically generated. You should know what you did if you 
want to edit this
+-- !merge_full_key --
+1      1       20      30
+1      2       35      40
+2      1       40      45
+2      2       45      \N
+3      1       \N      \N
+
+-- !merge_group_key --
+1      35      70
+2      45      45
+3      \N      \N
+
+-- !merge_global --
+45     115
+
+-- !merge_alias --
+1      35      70
+2      45      45
+3      \N      \N
+
+-- !merge_key_filter --
+1      20      30
+2      40      45
+3      \N      \N
+
+-- !union_phase_1 --
+1      35      70
+2      45      45
+3      \N      \N
+
+-- !merge_phase_1 --
+1      35      70
+2      45      45
+3      \N      \N
+
+-- !union_phase_2 --
+1      35      70
+2      45      45
+3      \N      \N
+
+-- !merge_phase_2 --
+1      35      70
+2      45      45
+3      \N      \N
+
+-- !union_global --
+45     115
+
+-- !empty_merge --
+\N     \N
+
+-- !empty_union --
+\N     \N
+
+-- !merge_with_count --
+1      35      70      2
+2      45      45      2
+3      \N      \N      1
+
+-- !union_with_count --
+1      35      70      2
+2      45      45      2
+3      \N      \N      1
+
diff --git 
a/regression-test/data/query_p0/outfile/agg_state/test_outfile_agg_state.out 
b/regression-test/data/query_p0/outfile/agg_state/test_outfile_agg_state.out
index a6c9b9dbf8e..77c2d4d529a 100644
--- a/regression-test/data/query_p0/outfile/agg_state/test_outfile_agg_state.out
+++ b/regression-test/data/query_p0/outfile/agg_state/test_outfile_agg_state.out
@@ -1,9 +1,9 @@
 -- This file is automatically generated. You should know what you did if you 
want to edit this
 -- !test --
-1      2       bb,a
+1      2       a,bb
 2      1       ccc
 
 -- !test --
-1      2       bb,a
+1      2       a,bb
 2      1       ccc
 
diff --git 
a/regression-test/data/query_p0/outfile/agg_state_array/test_outfile_agg_array.out
 
b/regression-test/data/query_p0/outfile/agg_state_array/test_outfile_agg_array.out
index fd97230395a..f694ddd6cef 100644
--- 
a/regression-test/data/query_p0/outfile/agg_state_array/test_outfile_agg_array.out
+++ 
b/regression-test/data/query_p0/outfile/agg_state_array/test_outfile_agg_array.out
@@ -1,6 +1,6 @@
 -- This file is automatically generated. You should know what you did if you 
want to edit this
 -- !test --
-1      [2, 1]
+1      [1, 2]
 2      [3]
 
 -- !test --
diff --git 
a/regression-test/suites/cloud_p0/different_serialize/different_serialize.groovy
 
b/regression-test/suites/cloud_p0/different_serialize/different_serialize.groovy
index df452fe8446..d11ffe36d2a 100644
--- 
a/regression-test/suites/cloud_p0/different_serialize/different_serialize.groovy
+++ 
b/regression-test/suites/cloud_p0/different_serialize/different_serialize.groovy
@@ -77,20 +77,20 @@ suite ("different_serialize_cloud") {
     qt_select_mv "select k1,bitmap_count(bitmap_agg(k2)) from d_table group by 
k1 order by 1;"
 
     explain {
-        sql("select k1,array_agg(k2) from d_table group by k1 order by 1;")
+        sql("select k1,array_sort(array_agg(k2)) from d_table group by k1 
order by 1;")
         contains "(mv3)"
     }
-    qt_select_mv "select k1,array_agg(k2) from d_table group by k1 order by 1;"
+    qt_select_mv "select k1,array_sort(array_agg(k2)) from d_table group by k1 
order by 1;"
 
     explain {
-        sql("select k1,collect_list(k2,3) from d_table group by k1 order by 
1;")
+        sql("select k1,array_sort(collect_list(k2,3)) from d_table group by k1 
order by 1;")
         contains "(mv4)"
     }
-    qt_select_mv "select k1,collect_list(k2,3) from d_table group by k1 order 
by 1;"
+    qt_select_mv "select k1,array_sort(collect_list(k2,3)) from d_table group 
by k1 order by 1;"
 
     explain {
-        sql("select k1,collect_set(k2,3) from d_table group by k1 order by 1;")
+        sql("select k1,array_sort(collect_set(k2,3)) from d_table group by k1 
order by 1;")
         contains "(mv5)"
     }
-    qt_select_mv "select k1,collect_set(k2,3) from d_table group by k1 order 
by 1;"
+    qt_select_mv "select k1,array_sort(collect_set(k2,3)) from d_table group 
by k1 order by 1;"
 }
diff --git a/regression-test/suites/datatype_p0/agg_state/array/array.groovy 
b/regression-test/suites/datatype_p0/agg_state/array/array.groovy
index 0af72638e53..9fd8b1ff7a1 100644
--- a/regression-test/suites/datatype_p0/agg_state/array/array.groovy
+++ b/regression-test/suites/datatype_p0/agg_state/array/array.groovy
@@ -37,5 +37,5 @@ suite("test_agg_state_array") {
     sql "insert into a_table values(1,array_agg_state(2));"
     sql "insert into a_table values(2,array_agg_state(3));"
 
-    qt_test "select k1,array_agg_merge(k2) from a_table group by k1 order by 
k1;"
+    qt_test "select k1,array_sort(array_agg_merge(k2)) from a_table group by 
k1 order by k1;"
 }
diff --git 
a/regression-test/suites/datatype_p0/agg_state/group_concat/test_agg_state_group_concat.groovy
 
b/regression-test/suites/datatype_p0/agg_state/group_concat/test_agg_state_group_concat.groovy
index f17e44a1596..f1f48d6d6e3 100644
--- 
a/regression-test/suites/datatype_p0/agg_state/group_concat/test_agg_state_group_concat.groovy
+++ 
b/regression-test/suites/datatype_p0/agg_state/group_concat/test_agg_state_group_concat.groovy
@@ -33,12 +33,12 @@ suite("test_agg_state_group_concat") {
     sql "insert into a_table values(1,group_concat_state('ccc'));"
 
     qt_length1 """select k1,length(k2) from a_table order by k1;"""
-    qt_group1 """select k1,group_concat_merge(k2) from a_table group by k1 
order by k1;"""
-    qt_merge1 """select group_concat_merge(k2) from a_table;"""
+    qt_group1 """select 
k1,array_join(array_sort(split_by_string(group_concat_merge(k2), ',')), ',') 
from a_table group by k1 order by k1;"""
+    qt_merge1 """select 
array_join(array_sort(split_by_string(group_concat_merge(k2), ',')), ',') from 
a_table;"""
     
     qt_length2 """select k1,length(k2) from a_table order by k1;"""
-    qt_group2 """select k1,group_concat_merge(k2) from a_table group by k1 
order by k1;"""
-    qt_merge2 """select group_concat_merge(k2) from a_table;"""
+    qt_group2 """select 
k1,array_join(array_sort(split_by_string(group_concat_merge(k2), ',')), ',') 
from a_table group by k1 order by k1;"""
+    qt_merge2 """select 
array_join(array_sort(split_by_string(group_concat_merge(k2), ',')), ',') from 
a_table;"""
     
-    qt_union """ select group_concat_merge(kstate) from (select 
k1,group_concat_union(k2) kstate from a_table group by k1 order by k1) t; """
+    qt_union """ select 
array_join(array_sort(split_by_string(group_concat_merge(kstate), ',')), ',') 
from (select k1,group_concat_union(k2) kstate from a_table group by k1 order by 
k1) t; """
 }
diff --git a/regression-test/suites/datatype_p0/agg_state/map/map.groovy 
b/regression-test/suites/datatype_p0/agg_state/map/map.groovy
index 8920ed70e57..5a0f82468ab 100644
--- a/regression-test/suites/datatype_p0/agg_state/map/map.groovy
+++ b/regression-test/suites/datatype_p0/agg_state/map/map.groovy
@@ -33,15 +33,26 @@ suite("test_agg_state_map") {
     distributed BY hash(k1) buckets 3
     properties("replication_num" = "1");
     """
+    // Unique map keys within each group make the values independent of state 
merge order.
     sql "insert into a_table values(1,map_agg_state(1,1));"
     sql "insert into a_table values(1,map_agg_state(2,2));"
-    sql "insert into a_table values(1,map_agg_state(1,11));"
-    sql "insert into a_table values(1,map_agg_state(2,22));"
+    sql "insert into a_table values(1,map_agg_state(5,11));"
+    sql "insert into a_table values(1,map_agg_state(6,22));"
     sql "insert into a_table values(1,map_agg_state(null, 100));"
     sql "insert into a_table values(2,map_agg_state(3,3));"
     sql "insert into a_table values(2,map_agg_state(4,null));"
     sql "insert into a_table values(2,map_agg_state(null,null));"
-    sql "insert into a_table values(2,map_agg_state(null,400));"
+    sql "insert into a_table values(2,map_agg_state(7,400));"
 
-    qt_test "select k1,map_agg_merge(k2) from a_table group by k1 order by k1;"
+    def query = """
+        select k1,
+               array_sort(map_keys(map_agg_merge(k2))),
+               array_sortby(map_values(map_agg_merge(k2)), 
map_keys(map_agg_merge(k2)))
+        from a_table group by k1 order by k1;
+    """
+    explain {
+        sql query
+        contains "(a_table), PREAGGREGATION: ON"
+    }
+    qt_test query
 }
diff --git 
a/regression-test/suites/mv_p0/agg_state/diffrent_serialize/diffrent_serialize.groovy
 
b/regression-test/suites/mv_p0/agg_state/diffrent_serialize/diffrent_serialize.groovy
index ea3d2467583..cd128efd959 100644
--- 
a/regression-test/suites/mv_p0/agg_state/diffrent_serialize/diffrent_serialize.groovy
+++ 
b/regression-test/suites/mv_p0/agg_state/diffrent_serialize/diffrent_serialize.groovy
@@ -67,22 +67,22 @@ suite ("diffrent_serialize") {
     }
     qt_select_mv "select 
k1,array_sort(map_keys(map_agg(k2,k3))),array_sortby(map_values(map_agg(k2,k3)),map_keys(map_agg(k2,k3)))
 from d_table group by k1 order by 1;"
     explain {
-        sql("select k1,array_agg(k2) from d_table group by k1 order by 1;")
+        sql("select k1,array_sort(array_agg(k2)) from d_table group by k1 
order by 1;")
         contains "(mv3)"
     }
-    qt_select_mv "select k1,array_agg(k2) from d_table group by k1 order by 1;"
+    qt_select_mv "select k1,array_sort(array_agg(k2)) from d_table group by k1 
order by 1;"
 
     explain {
-        sql("select k1,collect_list(k2,3) from d_table group by k1 order by 
1;")
+        sql("select k1,array_sort(collect_list(k2,3)) from d_table group by k1 
order by 1;")
         contains "(mv4)"
     }
-    qt_select_mv "select k1,collect_list(k2,3) from d_table group by k1 order 
by 1;"
+    qt_select_mv "select k1,array_sort(collect_list(k2,3)) from d_table group 
by k1 order by 1;"
 
     explain {
-        sql("select k1,collect_set(k2,3) from d_table group by k1 order by 1;")
+        sql("select k1,array_sort(collect_set(k2,3)) from d_table group by k1 
order by 1;")
         contains "(mv5)"
     }
-    qt_select_mv "select k1,collect_set(k2,3) from d_table group by k1 order 
by 1;"
+    qt_select_mv "select k1,array_sort(collect_set(k2,3)) from d_table group 
by k1 order by 1;"
 
     sql "insert into d_table select 1,1,1,'a';"
     sql "insert into d_table select 1,2,1,'a';"
@@ -93,7 +93,7 @@ suite ("diffrent_serialize") {
     mv_rewrite_success("select k1, multi_distinct_sum(k3) from d_table group 
by k1 order by k1;", "mv1_3")
     qt_select_mv "select k1, multi_distinct_sum(k3) from d_table group by k1 
order by k1;"
 
-    mv_rewrite_success("select k1, multi_distinct_group_concat(k4) from 
d_table group by k1 order by k1;", "mv1_2")
-    qt_select_mv "select k1, multi_distinct_group_concat(k4) from d_table 
group by k1 order by k1;"
+    mv_rewrite_success("select k1, 
array_join(array_sort(split_by_string(multi_distinct_group_concat(k4), ',')), 
',') from d_table group by k1 order by k1;", "mv1_2")
+    qt_select_mv "select k1, 
array_join(array_sort(split_by_string(multi_distinct_group_concat(k4), ',')), 
',') from d_table group by k1 order by k1;"
 
 }
diff --git 
a/regression-test/suites/nereids_rules_p0/set_preagg/agg_state_preagg.groovy 
b/regression-test/suites/nereids_rules_p0/set_preagg/agg_state_preagg.groovy
new file mode 100644
index 00000000000..a7190ffa5b0
--- /dev/null
+++ b/regression-test/suites/nereids_rules_p0/set_preagg/agg_state_preagg.groovy
@@ -0,0 +1,137 @@
+// 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("agg_state_preagg") {
+    sql "SET enable_agg_state = true"
+    sql "DROP TABLE IF EXISTS agg_state_preagg_generic"
+    sql """
+        CREATE TABLE agg_state_preagg_generic (
+            k1 INT,
+            k2 INT,
+            m AGG_STATE<max_by(INT NOT NULL, INT)> GENERIC,
+            s AGG_STATE<sum(INT)> GENERIC
+        )
+        AGGREGATE KEY(k1, k2)
+        DISTRIBUTED BY HASH(k2) BUCKETS 3
+        PROPERTIES("replication_num" = "1", "disable_auto_compaction" = "true")
+    """
+    // Separate rowsets keep partial states for the same full key. The SUM 
state
+    // also checks multiplicity, which an idempotent MAX_BY alone would not 
detect.
+    sql """
+        INSERT INTO agg_state_preagg_generic VALUES
+            (1, 1, max_by_state(10, 1), sum_state(10)),
+            (1, 2, max_by_state(30, 3), sum_state(NULL)),
+            (2, 1, max_by_state(40, 4), sum_state(40)),
+            (3, 1, max_by_state(50, NULL), sum_state(NULL))
+    """
+    sql """
+        INSERT INTO agg_state_preagg_generic VALUES
+            (1, 1, max_by_state(20, 2), sum_state(20)),
+            (1, 2, max_by_state(35, 5), sum_state(35)),
+            (2, 1, max_by_state(40, 4), sum_state(5)),
+            (3, 1, max_by_state(60, NULL), sum_state(NULL))
+    """
+    sql """
+        INSERT INTO agg_state_preagg_generic VALUES
+            (1, 1, max_by_state(5, NULL), sum_state(NULL)),
+            (1, 2, max_by_state(35, 5), sum_state(5)),
+            (2, 2, max_by_state(45, 6), sum_state(NULL))
+    """
+
+    def mergeQueries = [
+        full_key: """SELECT k1, k2, max_by_merge(m), sum_merge(s)
+                     FROM agg_state_preagg_generic GROUP BY k1, k2""",
+        group_key: """SELECT k1, max_by_merge(m), sum_merge(s)
+                      FROM agg_state_preagg_generic GROUP BY k1""",
+        global: """SELECT max_by_merge(m), sum_merge(s) FROM 
agg_state_preagg_generic""",
+        alias: """SELECT k1, max_by_merge(ms), sum_merge(ss)
+                  FROM (SELECT k1, m AS ms, s AS ss FROM 
agg_state_preagg_generic) t
+                  GROUP BY k1""",
+        key_filter: """SELECT k1, max_by_merge(m), sum_merge(s)
+                       FROM agg_state_preagg_generic WHERE k2 = 1 GROUP BY 
k1"""
+    ]
+    mergeQueries.each { name, query ->
+        explain {
+            sql query
+            contains "(agg_state_preagg_generic), PREAGGREGATION: ON"
+        }
+        "order_qt_merge_${name}"(query)
+    }
+
+    for (def phase : [1, 2]) {
+        def unionQuery = """
+            SELECT /*+ SET_VAR(agg_phase=${phase}) */ k1, max_by_union(m) m, 
sum_union(s) s
+            FROM agg_state_preagg_generic GROUP BY k1
+        """
+        explain {
+            sql unionQuery
+            contains "(agg_state_preagg_generic), PREAGGREGATION: ON"
+        }
+        "order_qt_union_phase_${phase}"("""
+            SELECT k1, max_by_merge(m), sum_merge(s) FROM (${unionQuery}) t 
GROUP BY k1
+        """)
+        "order_qt_merge_phase_${phase}"("""
+            SELECT /*+ SET_VAR(agg_phase=${phase}) */ k1, max_by_merge(m), 
sum_merge(s)
+            FROM agg_state_preagg_generic GROUP BY k1
+        """)
+    }
+    order_qt_union_global """
+        SELECT max_by_merge(m), sum_merge(s) FROM (
+            SELECT max_by_union(m) m, sum_union(s) s FROM 
agg_state_preagg_generic
+        ) t
+    """
+    order_qt_empty_merge """
+        SELECT max_by_merge(m), sum_merge(s) FROM agg_state_preagg_generic 
WHERE k1 = 100
+    """
+    order_qt_empty_union """
+        SELECT max_by_merge(m), sum_merge(s) FROM (
+            SELECT max_by_union(m) m, sum_union(s) s
+            FROM agg_state_preagg_generic WHERE k1 = 100
+        ) t
+    """
+
+    // A sibling COUNT needs storage-merged rows. Enabling merge/union must not
+    // bypass checks on other aggregates, value filters, or argument 
expressions.
+    for (def suffix : ["merge", "union"]) {
+        explain {
+            sql """SELECT k1, max_by_${suffix}(m), count(*)
+                   FROM agg_state_preagg_generic GROUP BY k1"""
+            contains "(agg_state_preagg_generic), PREAGGREGATION: OFF"
+        }
+        explain {
+            sql """SELECT k1, max_by_${suffix}(m)
+                   FROM agg_state_preagg_generic WHERE length(m) > 0 GROUP BY 
k1"""
+            contains "(agg_state_preagg_generic), PREAGGREGATION: OFF"
+        }
+        explain {
+            sql """SELECT k1, max_by_${suffix}(if(k2 = 1, m, NULL))
+                   FROM agg_state_preagg_generic GROUP BY k1"""
+            contains "(agg_state_preagg_generic), PREAGGREGATION: OFF"
+        }
+    }
+    order_qt_merge_with_count """
+        SELECT k1, max_by_merge(m), sum_merge(s), count(*)
+        FROM agg_state_preagg_generic GROUP BY k1
+    """
+    order_qt_union_with_count """
+        SELECT k1, max_by_merge(m), sum_merge(s), max(n) FROM (
+            SELECT k1, max_by_union(m) m, sum_union(s) s, count(*) n
+            FROM agg_state_preagg_generic GROUP BY k1
+        ) t GROUP BY k1
+    """
+
+}
diff --git 
a/regression-test/suites/query_p0/outfile/agg_state/test_outfile_agg_state.groovy
 
b/regression-test/suites/query_p0/outfile/agg_state/test_outfile_agg_state.groovy
index baa0c8adf5c..1e808246b68 100644
--- 
a/regression-test/suites/query_p0/outfile/agg_state/test_outfile_agg_state.groovy
+++ 
b/regression-test/suites/query_p0/outfile/agg_state/test_outfile_agg_state.groovy
@@ -45,7 +45,7 @@ suite("test_outfile_agg_state") {
     sql "insert into a_table 
values(1,max_by_state(2,2),group_concat_state('bb'));"
     sql "insert into a_table 
values(2,max_by_state(1,3),group_concat_state('ccc'));"
 
-    qt_test "select k1,max_by_merge(k2),group_concat_merge(k3) from a_table 
group by k1 order by k1;"
+    qt_test "select 
k1,max_by_merge(k2),array_join(array_sort(split_by_string(group_concat_merge(k3),
 ',')), ',') from a_table group by k1 order by k1;"
 
     sql """select * from a_table into outfile 
"file://${testHelper.remoteDir}/e_" FORMAT AS PARQUET;"""
     testHelper.collect()
@@ -67,7 +67,7 @@ suite("test_outfile_agg_state") {
     curl --location-trusted -u 
${context.config.jdbcUser}:${context.config.jdbcPassword} -H "format:PARQUET" 
-H "Expect:100-continue" -T ${filePath} -XPUT 
${getDorisHttpScheme()}://${context.config.feHttpAddress}/api/regression_test_query_p0_outfile_agg_state/a_table2/_stream_load${getDorisCurlTlsOptions()}
     """
     Thread.sleep(10000)
-    qt_test "select k1,max_by_merge(k2),group_concat_merge(k3) from a_table2 
group by k1 order by k1;"
+    qt_test "select 
k1,max_by_merge(k2),array_join(array_sort(split_by_string(group_concat_merge(k3),
 ',')), ',') from a_table2 group by k1 order by k1;"
 
     testHelper.close()
 }
diff --git 
a/regression-test/suites/query_p0/outfile/agg_state_array/test_outfile_agg_array.groovy
 
b/regression-test/suites/query_p0/outfile/agg_state_array/test_outfile_agg_array.groovy
index af9c3ad6f76..6428e40fd07 100644
--- 
a/regression-test/suites/query_p0/outfile/agg_state_array/test_outfile_agg_array.groovy
+++ 
b/regression-test/suites/query_p0/outfile/agg_state_array/test_outfile_agg_array.groovy
@@ -45,7 +45,7 @@ suite("test_outfile_agg_state_array") {
     sql "insert into a_table values(1,array_agg_state(2));"
     sql "insert into a_table values(2,array_agg_state(3));"
 
-    qt_test "select k1,array_agg_merge(k2) from a_table group by k1 order by 
k1;"
+    qt_test "select k1,array_sort(array_agg_merge(k2)) from a_table group by 
k1 order by k1;"
 
     sql """select * from a_table into outfile 
"file://${testHelper.remoteDir}/tmp_" FORMAT AS PARQUET;"""
     testHelper.collect()
@@ -67,6 +67,6 @@ suite("test_outfile_agg_state_array") {
     curl --location-trusted -u 
${context.config.jdbcUser}:${context.config.jdbcPassword} -H "format:PARQUET" 
-H "Expect:100-continue" -T ${filePath} 
${getDorisHttpScheme()}://${context.config.feHttpAddress}/api/regression_test_query_p0_outfile_agg_state_array/a_table2/_stream_load${getDorisCurlTlsOptions()}
     """
     Thread.sleep(10000)
-    qt_test "select k1,max_by_merge(k2),group_concat_merge(k3) from a_table2 
group by k1 order by k1;"
+    qt_test "select 
k1,max_by_merge(k2),array_join(array_sort(split_by_string(group_concat_merge(k3),
 ',')), ',') from a_table2 group by k1 order by k1;"
     testHelper.close()
 }


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

Reply via email to