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

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


The following commit(s) were added to refs/heads/master by this push:
     new b62d37017e4 [FLINK-40214][table] Preserve JSON_OBJECTAGG key when 
wrapping value (#28795)
b62d37017e4 is described below

commit b62d37017e4cd507f40492a805dbcf554a98ffb2
Author: Au-Miner <[email protected]>
AuthorDate: Thu Jul 30 17:24:22 2026 +0800

    [FLINK-40214][table] Preserve JSON_OBJECTAGG key when wrapping value 
(#28795)
---
 .../logical/WrapJsonAggFunctionArgumentsRule.java  | 31 +++++++++++-----------
 .../functions/JsonAggregationFunctionsITCase.java  | 30 +++++++++++++++++++++
 .../WrapJsonAggFunctionArgumentsRuleTest.xml       | 20 +++++++-------
 3 files changed, 55 insertions(+), 26 deletions(-)

diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/rules/logical/WrapJsonAggFunctionArgumentsRule.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/rules/logical/WrapJsonAggFunctionArgumentsRule.java
index 089e542f4c9..100141d9397 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/rules/logical/WrapJsonAggFunctionArgumentsRule.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/rules/logical/WrapJsonAggFunctionArgumentsRule.java
@@ -38,9 +38,6 @@ import org.apache.calcite.sql.SqlAggFunction;
 import org.apache.calcite.sql.fun.SqlJsonArrayAggAggFunction;
 import org.apache.calcite.sql.fun.SqlJsonObjectAggAggFunction;
 import org.apache.calcite.tools.RelBuilder;
-import org.apache.calcite.util.mapping.MappingType;
-import org.apache.calcite.util.mapping.Mappings;
-import org.apache.calcite.util.mapping.Mappings.TargetMapping;
 import org.immutables.value.Value;
 
 import java.util.ArrayList;
@@ -106,12 +103,9 @@ public class WrapJsonAggFunctionArgumentsRule
         Map<Integer, Integer> wrapIndicesMap = new HashMap<>();
         for (int i = 0; i < aggCallList.size(); i++) {
             AggregateCall currentCall = aggCallList.get(i);
-            if (currentCall.getAggregation() instanceof 
SqlJsonObjectAggAggFunction) {
-                // For JSON_OBJECTAGG we only need to wrap its second (= 
value) argument
-                final int valueIndex = currentCall.getArgList().get(1);
-                wrapIndicesMap.put(i, valueIndex);
-            } else if (currentCall.getAggregation() instanceof 
SqlJsonArrayAggAggFunction) {
-                final int valueIndex = currentCall.getArgList().get(0);
+            if (isJsonAggregation(currentCall)) {
+                final int valueIndex =
+                        
currentCall.getArgList().get(getValueArgPosition(currentCall));
                 wrapIndicesMap.put(i, valueIndex);
             }
         }
@@ -126,14 +120,13 @@ public class WrapJsonAggFunctionArgumentsRule
                 valueIndicesAfterProjection);
 
         List<AggregateCall> newWrappedArgCallList = new 
ArrayList<>(aggCallList);
-        final int newInputCount = inputCount + 
valueIndicesAfterProjection.size();
         for (Integer jsonAggCallIndex : wrapIndicesMap.keySet()) {
-            final TargetMapping argsMapping =
-                    Mappings.create(MappingType.BIJECTION, newInputCount, 
newInputCount);
-            Integer valueIndex = wrapIndicesMap.get(jsonAggCallIndex);
-            argsMapping.set(valueIndex, 
valueIndicesAfterProjection.get(valueIndex));
-            final AggregateCall newAggregateCall =
-                    
newWrappedArgCallList.get(jsonAggCallIndex).transform(argsMapping);
+            final AggregateCall aggregateCall = 
newWrappedArgCallList.get(jsonAggCallIndex);
+            final List<Integer> newArgList = new 
ArrayList<>(aggregateCall.getArgList());
+            final int valueArgPosition = getValueArgPosition(aggregateCall);
+            final Integer valueIndex = wrapIndicesMap.get(jsonAggCallIndex);
+            newArgList.set(valueArgPosition, 
valueIndicesAfterProjection.get(valueIndex));
+            final AggregateCall newAggregateCall = 
aggregateCall.withArgList(newArgList);
             newWrappedArgCallList.set(jsonAggCallIndex, newAggregateCall);
         }
 
@@ -178,6 +171,12 @@ public class WrapJsonAggFunctionArgumentsRule
                 || aggregation instanceof SqlJsonArrayAggAggFunction;
     }
 
+    private static int getValueArgPosition(AggregateCall jsonAggCall) {
+        // AggregateCall argument positions are zero-based: JSON_OBJECTAGG has 
(key, value),
+        // whereas JSON_ARRAYAGG has only (value).
+        return jsonAggCall.getAggregation() instanceof 
SqlJsonObjectAggAggFunction ? 1 : 0;
+    }
+
     // 
---------------------------------------------------------------------------------------------
 
     /** Configuration for {@link WrapJsonAggFunctionArgumentsRule}. */
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonAggregationFunctionsITCase.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonAggregationFunctionsITCase.java
index 32827b79d5a..4c352943b8d 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonAggregationFunctionsITCase.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonAggregationFunctionsITCase.java
@@ -60,6 +60,15 @@ class JsonAggregationFunctionsITCase extends 
BuiltInAggregateFunctionTestBase {
                                 ROW(VARCHAR(2000).notNull()),
                                 ROW(STRING().notNull()),
                                 
Collections.singletonList(Row.of("{\"A\":1,\"B\":null,\"C\":3}"))),
+                
TestSpec.forFunction(BuiltInFunctionDefinitions.JSON_OBJECTAGG_NULL_ON_NULL)
+                        .withDescription("Aggregation With Same Key And Value")
+                        .withSource(
+                                ROW(STRING()),
+                                Arrays.asList(Row.ofKind(INSERT, "A"), 
Row.ofKind(INSERT, "B")))
+                        .testSqlResult(
+                                source -> "SELECT JSON_OBJECTAGG(f0 VALUE f0) 
FROM " + source,
+                                ROW(VARCHAR(2000).notNull()),
+                                
Collections.singletonList(Row.of("{\"A\":\"A\",\"B\":\"B\"}"))),
                 
TestSpec.forFunction(BuiltInFunctionDefinitions.JSON_OBJECTAGG_ABSENT_ON_NULL)
                         .withDescription("Omits NULLs")
                         .withSource(
@@ -273,6 +282,27 @@ class JsonAggregationFunctionsITCase extends 
BuiltInAggregateFunctionTestBase {
                                                 + " GROUP BY TUMBLE(f2, 
INTERVAL '5' SECOND)",
                                 ROW(VARCHAR(2000).notNull()),
                                 Arrays.asList(Row.of("{\"A\":1,\"B\":2}"), 
Row.of("{\"C\":3}"))),
+                
TestSpec.forFunction(BuiltInFunctionDefinitions.JSON_OBJECTAGG_NULL_ON_NULL)
+                        .withDescription("Window Aggregation With Same Key And 
Value")
+                        .withSource(
+                                ROW(STRING(), TIMESTAMP(3)),
+                                Arrays.asList(
+                                        Row.ofKind(
+                                                INSERT,
+                                                "A",
+                                                
LocalDateTime.parse("2020-01-01T00:00:01")),
+                                        Row.ofKind(
+                                                INSERT,
+                                                "B",
+                                                
LocalDateTime.parse("2020-01-01T00:00:02"))))
+                        .withWatermark("f1", "f1 - INTERVAL '1' SECOND")
+                        .testSqlResult(
+                                source ->
+                                        "SELECT JSON_OBJECTAGG(f0 VALUE f0) 
FROM "
+                                                + source
+                                                + " GROUP BY TUMBLE(f1, 
INTERVAL '5' SECOND)",
+                                ROW(VARCHAR(2000).notNull()),
+                                
Collections.singletonList(Row.of("{\"A\":\"A\",\"B\":\"B\"}"))),
                 
TestSpec.forFunction(BuiltInFunctionDefinitions.JSON_OBJECTAGG_NULL_ON_NULL)
                         .withDescription("Window Group Aggregation With Other 
Aggs")
                         .withSource(
diff --git 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/logical/WrapJsonAggFunctionArgumentsRuleTest.xml
 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/logical/WrapJsonAggFunctionArgumentsRuleTest.xml
index e5480a1e27b..5f1a839794a 100644
--- 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/logical/WrapJsonAggFunctionArgumentsRuleTest.xml
+++ 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/rules/logical/WrapJsonAggFunctionArgumentsRuleTest.xml
@@ -718,7 +718,7 @@ LogicalAggregate(group=[{}], EXPR$0=[COUNT()], 
EXPR$1=[JSON_OBJECTAGG_NULL_ON_NU
     </Resource>
     <Resource name="optimized rel plan">
       <![CDATA[
-GroupAggregate(select=[COUNT(*) AS EXPR$0, JSON_OBJECTAGG_NULL_ON_NULL($f1, 
$f1) AS EXPR$1])
+GroupAggregate(select=[COUNT(*) AS EXPR$0, JSON_OBJECTAGG_NULL_ON_NULL(f1, 
$f1) AS EXPR$1])
 +- Exchange(distribution=[single])
    +- Calc(select=[f1, JSON_STRING(f1) AS $f1])
       +- TableSourceScan(table=[[default_catalog, default_database, T, 
project=[f1], metadata=[]]], fields=[f1])
@@ -740,7 +740,7 @@ LogicalProject(EXPR$0=[$1], EXPR$1=[$2])
     </Resource>
     <Resource name="optimized rel plan">
       <![CDATA[
-GroupWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[COUNT(*) AS EXPR$0, JSON_OBJECTAGG_NULL_ON_NULL($f2, $f2) AS EXPR$1])
+GroupWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[COUNT(*) AS EXPR$0, JSON_OBJECTAGG_NULL_ON_NULL(f1, $f2) AS EXPR$1])
 +- Exchange(distribution=[single])
    +- Calc(select=[rt, f1, JSON_STRING(f1) AS $f2])
       +- WatermarkAssigner(rowtime=[rt], watermark=[-(rt, 1000:INTERVAL 
SECOND)])
@@ -763,7 +763,7 @@ LogicalAggregate(group=[{}], EXPR$0=[COUNT()], 
EXPR$1=[JSON_OBJECTAGG_NULL_ON_NU
       <![CDATA[
 SortAggregate(isMerge=[true], select=[Final_COUNT(count1$0) AS EXPR$0, 
Final_JSON_OBJECTAGG_NULL_ON_NULL(EXPR$1) AS EXPR$1])
 +- Exchange(distribution=[single])
-   +- LocalSortAggregate(select=[Partial_COUNT(*) AS count1$0, 
Partial_JSON_OBJECTAGG_NULL_ON_NULL($f1, $f1) AS EXPR$1])
+   +- LocalSortAggregate(select=[Partial_COUNT(*) AS count1$0, 
Partial_JSON_OBJECTAGG_NULL_ON_NULL(f1, $f1) AS EXPR$1])
       +- Calc(select=[f1, JSON_STRING(f1) AS $f1])
          +- TableSourceScan(table=[[default_catalog, default_database, T, 
project=[f1], metadata=[]]], fields=[f1])
 ]]>
@@ -786,7 +786,7 @@ LogicalProject(EXPR$0=[$1], EXPR$1=[$2])
 SortWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[Final_COUNT(count1$0) AS EXPR$0, 
Final_JSON_OBJECTAGG_NULL_ON_NULL(EXPR$1) AS EXPR$1])
 +- Sort(orderBy=[assignedWindow$ ASC])
    +- Exchange(distribution=[single])
-      +- LocalSortWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[Partial_COUNT(*) AS count1$0, Partial_JSON_OBJECTAGG_NULL_ON_NULL($f2, 
$f2) AS EXPR$1])
+      +- LocalSortWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[Partial_COUNT(*) AS count1$0, Partial_JSON_OBJECTAGG_NULL_ON_NULL(f1, 
$f2) AS EXPR$1])
          +- Calc(select=[rt, f1, JSON_STRING(f1) AS $f2])
             +- Sort(orderBy=[rt ASC])
                +- TableSourceScan(table=[[default_catalog, default_database, 
T, project=[rt, f1], metadata=[]]], fields=[rt, f1])
@@ -809,7 +809,7 @@ LogicalAggregate(group=[{}], 
EXPR$0=[JSON_OBJECTAGG_NULL_ON_NULL($0, $0)])
     </Resource>
     <Resource name="optimized rel plan">
       <![CDATA[
-GroupAggregate(select=[JSON_OBJECTAGG_NULL_ON_NULL($f1, $f1) AS EXPR$0])
+GroupAggregate(select=[JSON_OBJECTAGG_NULL_ON_NULL(f1, $f1) AS EXPR$0])
 +- Exchange(distribution=[single])
    +- Calc(select=[f1, JSON_STRING(f1) AS $f1])
       +- WindowTableFunction(window=[TUMBLE(time_col=[rt], size=[5 s])])
@@ -835,7 +835,7 @@ LogicalAggregate(group=[{}], 
EXPR$0=[JSON_OBJECTAGG_NULL_ON_NULL($0, $0)])
       <![CDATA[
 SortAggregate(isMerge=[true], 
select=[Final_JSON_OBJECTAGG_NULL_ON_NULL(EXPR$0) AS EXPR$0])
 +- Exchange(distribution=[single])
-   +- LocalSortAggregate(select=[Partial_JSON_OBJECTAGG_NULL_ON_NULL($f1, $f1) 
AS EXPR$0])
+   +- LocalSortAggregate(select=[Partial_JSON_OBJECTAGG_NULL_ON_NULL(f1, $f1) 
AS EXPR$0])
       +- Calc(select=[f1, JSON_STRING(f1) AS $f1])
          +- WindowTableFunction(window=[TUMBLE(time_col=[rt], size=[5 s])])
             +- TableSourceScan(table=[[default_catalog, default_database, T, 
project=[f1, rt], metadata=[]]], fields=[f1, rt])
@@ -855,7 +855,7 @@ LogicalAggregate(group=[{}], 
EXPR$0=[JSON_OBJECTAGG_NULL_ON_NULL($0, $0)])
     </Resource>
     <Resource name="optimized rel plan">
       <![CDATA[
-GroupAggregate(select=[JSON_OBJECTAGG_NULL_ON_NULL($f1, $f1) AS EXPR$0])
+GroupAggregate(select=[JSON_OBJECTAGG_NULL_ON_NULL(f1, $f1) AS EXPR$0])
 +- Exchange(distribution=[single])
    +- Calc(select=[f1, JSON_STRING(f1) AS $f1])
       +- TableSourceScan(table=[[default_catalog, default_database, T, 
project=[f1], metadata=[]]], fields=[f1])
@@ -877,7 +877,7 @@ LogicalProject(EXPR$0=[$1])
     </Resource>
     <Resource name="optimized rel plan">
       <![CDATA[
-GroupWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[JSON_OBJECTAGG_NULL_ON_NULL($f2, $f2) AS EXPR$0])
+GroupWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[JSON_OBJECTAGG_NULL_ON_NULL(f1, $f2) AS EXPR$0])
 +- Exchange(distribution=[single])
    +- Calc(select=[rt, f1, JSON_STRING(f1) AS $f2])
       +- WatermarkAssigner(rowtime=[rt], watermark=[-(rt, 1000:INTERVAL 
SECOND)])
@@ -900,7 +900,7 @@ LogicalAggregate(group=[{}], 
EXPR$0=[JSON_OBJECTAGG_NULL_ON_NULL($0, $0)])
       <![CDATA[
 SortAggregate(isMerge=[true], 
select=[Final_JSON_OBJECTAGG_NULL_ON_NULL(EXPR$0) AS EXPR$0])
 +- Exchange(distribution=[single])
-   +- LocalSortAggregate(select=[Partial_JSON_OBJECTAGG_NULL_ON_NULL($f1, $f1) 
AS EXPR$0])
+   +- LocalSortAggregate(select=[Partial_JSON_OBJECTAGG_NULL_ON_NULL(f1, $f1) 
AS EXPR$0])
       +- Calc(select=[f1, JSON_STRING(f1) AS $f1])
          +- TableSourceScan(table=[[default_catalog, default_database, T, 
project=[f1], metadata=[]]], fields=[f1])
 ]]>
@@ -923,7 +923,7 @@ LogicalProject(EXPR$0=[$1])
 SortWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[Final_JSON_OBJECTAGG_NULL_ON_NULL(EXPR$0) AS EXPR$0])
 +- Sort(orderBy=[assignedWindow$ ASC])
    +- Exchange(distribution=[single])
-      +- LocalSortWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[Partial_JSON_OBJECTAGG_NULL_ON_NULL($f2, $f2) AS EXPR$0])
+      +- LocalSortWindowAggregate(window=[TumblingGroupWindow('w$, rt, 5000)], 
select=[Partial_JSON_OBJECTAGG_NULL_ON_NULL(f1, $f2) AS EXPR$0])
          +- Calc(select=[rt, f1, JSON_STRING(f1) AS $f2])
             +- Sort(orderBy=[rt ASC])
                +- TableSourceScan(table=[[default_catalog, default_database, 
T, project=[rt, f1], metadata=[]]], fields=[rt, f1])

Reply via email to