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]