SEPURI-SAI-KRISHNA commented on code in PR #29412:
URL: https://github.com/apache/flink/pull/29412#discussion_r4230584235


##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/runtime/stream/sql/NonTimeOverAggregateITCase.java:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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.
+ */
+
+package org.apache.flink.table.planner.runtime.stream.sql;
+
+import org.apache.flink.table.planner.factories.TestValuesTableFactory;
+import org.apache.flink.table.planner.runtime.utils.StreamingWithStateTestBase;
+import 
org.apache.flink.testutils.junit.extensions.parameterized.ParameterizedTestExtension;
+import org.apache.flink.testutils.junit.extensions.parameterized.Parameters;
+import org.apache.flink.types.Row;
+import org.apache.flink.types.RowKind;
+import org.apache.flink.util.CloseableIterator;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.TestTemplate;
+import org.junit.jupiter.api.extension.ExtendWith;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collection;
+import java.util.List;
+import java.util.Map;
+import java.util.TreeMap;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.entry;
+
+/** Tests for OVER windows ordered by a non-time attribute. */
+@ExtendWith(ParameterizedTestExtension.class)
+public class NonTimeOverAggregateITCase extends StreamingWithStateTestBase {
+
+    public NonTimeOverAggregateITCase(StateBackendMode state) {
+        super(state);
+    }
+
+    @Parameters(name = "backend = {0}")
+    public static Collection<Object[]> parameters() {
+        return Arrays.asList(
+                new Object[][] {
+                    {StreamingWithStateTestBase.HEAP_BACKEND()},
+                    {StreamingWithStateTestBase.ROCKSDB_BACKEND()}
+                });
+    }
+
+    @BeforeEach
+    public void before() {
+        super.before();
+        String dataId =
+                TestValuesTableFactory.registerData(
+                        Arrays.asList(
+                                Row.of("a", 30, "p"), Row.of("a", 20, "q"), 
Row.of("a", 10, "r")));
+        tEnv().executeSql(
+                        "CREATE TABLE src (k STRING, ord INT, v STRING) WITH ("
+                                + " 'connector' = 'values',"
+                                + " 'data-id' = '"
+                                + dataId
+                                + "',"
+                                + " 'bounded' = 'true')");
+        String repeatedDataId =
+                TestValuesTableFactory.registerData(
+                        Arrays.asList(
+                                Row.of("a", 30, "p"),
+                                Row.of("a", 10, "r"),
+                                Row.of("a", 20, "q"),
+                                Row.of("a", 25, "r")));
+        tEnv().executeSql(
+                        "CREATE TABLE src_repeated (k STRING, ord INT, v 
STRING) WITH ("
+                                + " 'connector' = 'values',"
+                                + " 'data-id' = '"
+                                + repeatedDataId
+                                + "',"
+                                + " 'bounded' = 'true')");
+        String appendDataId =
+                TestValuesTableFactory.registerData(
+                        Arrays.asList(Row.of("a", 10, "x"), Row.of("a", 20, 
"x")));
+        tEnv().executeSql(
+                        "CREATE TABLE src_append (k STRING, ord INT, v STRING) 
WITH ("
+                                + " 'connector' = 'values',"
+                                + " 'data-id' = '"
+                                + appendDataId
+                                + "',"
+                                + " 'bounded' = 'true')");
+    }
+
+    @TestTemplate
+    void testLag() throws Exception {
+        assertThat(currentValues("LAG(v, 1)"))
+                .containsExactly(entry(10, "null"), entry(20, "r"), entry(30, 
"q"));
+    }
+
+    @TestTemplate
+    void testLagInsertsInBetween() throws Exception {
+        assertThat(currentValues("src_repeated", "LAG(v, 1)"))
+                .containsExactly(entry(10, "null"), entry(20, "r"), entry(25, 
"q"), entry(30, "r"));
+    }
+
+    @TestTemplate
+    void testArrayAgg() throws Exception {
+        assertThat(currentValues("ARRAY_AGG(v)"))
+                .containsExactly(entry(10, "[r]"), entry(20, "[r, q]"), 
entry(30, "[r, q, p]"));
+    }
+
+    @TestTemplate
+    void testArrayAggInsertsInBetween() throws Exception {
+        assertThat(currentValues("src_repeated", "ARRAY_AGG(v)"))
+                .containsExactly(
+                        entry(10, "[r]"),
+                        entry(20, "[r, q]"),
+                        entry(25, "[r, q, r]"),
+                        entry(30, "[r, q, r, p]"));
+    }
+
+    @TestTemplate
+    void testBitmapBuildCardinalityInsertsInBetween() throws Exception {
+        assertThat(currentValues("src_repeated", 
"BITMAP_BUILD_CARDINALITY_AGG(ord)"))
+                .containsExactly(entry(10, "1"), entry(20, "2"), entry(25, 
"3"), entry(30, "4"));
+    }
+
+    @TestTemplate
+    void testCollect() throws Exception {
+        assertThat(currentValues("COLLECT(v)"))
+                .containsExactly(
+                        entry(10, "{r=1}"), entry(20, "{q=1, r=1}"), entry(30, 
"{p=1, q=1, r=1}"));
+    }
+
+    @TestTemplate
+    void testCollectCountsRepeatedValues() throws Exception {
+        assertThat(currentValues("src_repeated", "COLLECT(v)"))
+                .containsExactly(
+                        entry(10, "{r=1}"),
+                        entry(20, "{q=1, r=1}"),
+                        entry(25, "{q=1, r=2}"),
+                        entry(30, "{p=1, q=1, r=2}"));
+    }
+
+    @TestTemplate
+    void testCollectRows() throws Exception {
+        assertThat(currentValues("src_repeated", "COLLECT(v)", "ROWS UNBOUNDED 
PRECEDING"))
+                .containsExactly(
+                        entry(10, "{r=1}"),
+                        entry(20, "{q=1, r=1}"),
+                        entry(25, "{q=1, r=2}"),
+                        entry(30, "{p=1, q=1, r=2}"));
+    }
+
+    @TestTemplate
+    void testCollectAppendsAfterMaximum() throws Exception {
+        assertThat(currentValues("src_append", "COLLECT(v)"))
+                .containsExactly(entry(10, "{x=1}"), entry(20, "{x=2}"));
+    }
+
+    @TestTemplate
+    void testPercentile() throws Exception {
+        assertThat(currentValues("PERCENTILE(ord, 0.5)"))
+                .containsExactly(entry(10, "10.0"), entry(20, "15.0"), 
entry(30, "20.0"));
+    }
+
+    @TestTemplate
+    void testCountDistinct() throws Exception {
+        assertThat(currentValues("src_repeated", "COUNT(DISTINCT v)"))
+                .containsExactly(entry(10, "1"), entry(20, "2"), entry(25, 
"2"), entry(30, "3"));
+    }
+
+    @TestTemplate
+    void testCountDistinctAppendsAfterMaximum() throws Exception {
+        assertThat(currentValues("src_append", "COUNT(DISTINCT v)"))
+                .containsExactly(entry(10, "1"), entry(20, "1"));
+    }
+
+    @TestTemplate
+    void testRetractionsOfInsertInBetween() throws Exception {
+        assertThat(changelog("src_repeated", "ARRAY_AGG(v)"))
+                .containsExactly(
+                        "+I 30 [p]",
+                        "+I 10 [r]",
+                        "-U 30 [p]",
+                        "+U 30 [r, p]",
+                        "+I 20 [r, q]",
+                        "-U 30 [r, p]",
+                        "+U 30 [r, q, p]",
+                        "+I 25 [r, q, r]",
+                        "-U 30 [r, q, p]",
+                        "+U 30 [r, q, r, p]");
+    }
+
+    private static String overSql(String table, String agg, String frame) {
+        return "SELECT ord, "
+                + agg
+                + " OVER (PARTITION BY k ORDER BY ord "
+                + frame
+                + ") FROM "
+                + table;
+    }
+
+    /** Returns the raw changelog, so that retracted values are visible too. */
+    private List<String> changelog(String table, String agg) throws Exception {
+        final String frame = "";

Review Comment:
   Done.
   



##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecOverAggregate.java:
##########
@@ -321,7 +321,11 @@ private KeyedProcessFunction<RowData, RowData, RowData> 
createUnboundedOverProce
                         JavaScalaConversionUtil.toScala(aggCalls),
                         new boolean[aggCalls.size()],
                         false, // needInputCount
-                        true, // isStateBackendDataViews
+                        // The non-time functions keep one accumulator per 
sort key in state.
+                        // State backed data views are bound to the key, not 
to the sort key, so
+                        // all of those accumulators would share a single 
view. Keep the views in
+                        // the accumulator instead, so each sort key gets its 
own copy.
+                        timeAttribute != TimeAttribute.NON_TIME, // 
isStateBackendDataViews

Review Comment:
   You were right, they were starved behind the hung one. All six generated 
alone in about 5s each and are in, so OverAggregateRestoreTest goes from 10 to 
16 cases. The non-partitioned out-of-order one does fail alone, so five stay 
ignored: the four FLINK-38927 broke, plus that one. My earlier read was wrong, 
I took a repeated batch result as a property of the programs.
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to