MartijnVisser commented on code in PR #29412:
URL: https://github.com/apache/flink/pull/29412#discussion_r4205091522
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/over/AbstractNonTimeUnboundedPrecedingOver.java:
##########
@@ -418,10 +445,24 @@ Tuple2<Integer, Boolean> findIndexOfSortKey(
* @throws Exception
*/
RowData setAccumulatorAndGetValue(RowData accumulator) throws Exception {
- aggFuncs.setAccumulators(accumulator);
+ setAccumulators(accumulator);
return aggFuncs.getValue();
}
+ /**
+ * Hands an accumulator to the aggregate functions, copying it first if it
can hold a data view.
+ *
+ * <p>Accumulators are kept per sort key in {@link #accMapState}.
Accumulating into one that was
+ * read from state would also change the entry that state still holds for
an earlier sort key,
+ * because the heap state backend returns the stored object rather than a
copy.
+ *
+ * @param accumulator the accumulator to continue from
+ */
+ void setAccumulators(RowData accumulator) throws Exception {
Review Comment:
On heap, `(10, a), (20, a)` gives COUNT(DISTINCT v) 2 for 20. Nothing calls
this after the last `accMapState.put`, so the reset at the end of
`processElement` clears the stored distinct view.
##########
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:
This changes the accumulator state layout for aggregates with a data view.
FLINK-38765 bumped the sink node for that; I'd do the same here, plus a
non-time COLLECT and COUNT(DISTINCT) restore program.
##########
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:
On version 1 plans, a WARN when a non-time window uses a data view aggregate
would tell those users to recompile, since their results stay wrong. WDYT?
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/runtime/stream/sql/NonTimeOverAggregateITCase.java:
##########
@@ -0,0 +1,154 @@
+/*
+ * 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.Arrays;
+import java.util.Collection;
+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')");
+ }
+
+ @TestTemplate
+ void testLag() throws Exception {
+ assertThat(currentValues("LAG(v, 1)"))
+ .containsExactly(entry(10, "null"), entry(20, "r"), entry(30,
"q"));
+ }
+
+ @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 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 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 {
Review Comment:
`src_repeated` never appends after the current maximum, which is why heap
passes. Can you add a case like `(10, a), (20, a)`?
##########
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()],
Review Comment:
No call here needs retraction, so the retract variants of FIRST_VALUE,
LAST_VALUE and LISTAGG never get here, and JSON_OBJECTAGG doesn't parse with
OVER. Can you drop them from the release note and description?
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/over/AbstractNonTimeUnboundedPrecedingOver.java:
##########
@@ -231,6 +246,18 @@ public void open(OpenContext openContext) throws Exception
{
accMapState = getRuntimeContext().getMapState(accStateDescriptor);
+ // A data view is stored as a RAW field inside the accumulator. The
aggregate functions
+ // change such a view in place, and the heap state backend hands back
the very object it
+ // stores, so the accumulator has to be copied before it is
accumulated into. Accumulators
+ // without a RAW field are immutable from the operator's point of view
and need no copy.
+ if (Arrays.stream(accTypes)
Review Comment:
I'd drop this for the get/put copy suggested on FLINK-40735. That also
covers the reset, and the direct `setAccumulators` calls in FLINK-40883 that
make heap ROWS COLLECT retract unemitted rows.
--
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]