[
https://issues.apache.org/jira/browse/FLINK-40737?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Martijn Visser updated FLINK-40737:
-----------------------------------
Release Note:
COLLECT over an OVER window ordered by a non-time attribute returned a multiset
holding only the current row, PERCENTILE was wrong, and DISTINCT aggregates
counted values they had already seen, because all sort keys of a key shared one
state-backed data view. The views now live inside the per-sort-key accumulator
and these aggregates return correct results.
This changes the accumulator state layout of such aggregates, so
stream-exec-over-aggregate gained version 2 (minPlanVersion and minStateVersion
2.4). A compiled plan on version 1 still restores but keeps returning wrong
results for these aggregates; the planner logs a warning naming them. A
recompiled plan cannot be restored from a savepoint taken with the old one,
because the operator keeps its uid and the accumulator serializer is
incompatible; such jobs have to start fresh. Jobs without a compiled plan fail
at restore with a serializer incompatibility and also have to start fresh.
was:
COLLECT over an OVER window ordered by a non-time attribute returned a multiset
with only the
current row, PERCENTILE was wrong, and a DISTINCT aggregate counted values it
had already seen. On
the heap state backend the window also corrupted accumulators that state held
for another sort key,
which affected ARRAY_AGG, LAG, the BITMAP aggregates and user defined aggregate
functions whose
accumulator holds an array or a byte[]. All of these now return correct results.
This changes the accumulator state layout of aggregates that keep a data view,
so the ExecNode
stream-exec-over-aggregate gained a version 2. A plan still on version 1
restores as before but
keeps returning wrong results, and the planner logs a warning naming the
aggregates. Recompile the
plan to fix them; the operator then starts without its state.
> COLLECT over a non-time OVER window returns only the current row instead of
> the running window
> ----------------------------------------------------------------------------------------------
>
> Key: FLINK-40737
> URL: https://issues.apache.org/jira/browse/FLINK-40737
> Project: Flink
> Issue Type: Bug
> Components: Table SQL / Runtime
> Affects Versions: 2.3.0, 2.2.1, 2.1.3
> Reporter: sepuri sai krishna
> Assignee: sepuri sai krishna
> Priority: Major
> Labels: pull-request-available
> Attachments: CollectNonTimeOverRepro.java, pom.xml
>
>
> In streaming mode, {{COLLECT}} over an {{OVER}} window ordered by a non-time
> attribute
> returns a multiset containing only the current row, instead of the running
> window. No error
> is raised.
> {code:sql}
> SELECT ord, COLLECT(v) OVER (PARTITION BY k ORDER BY ord),
> ARRAY_AGG(v) OVER (PARTITION BY k ORDER BY ord),
> COUNT(*) OVER (PARTITION BY k ORDER BY ord)
> FROM (VALUES ('a',10,'p'),('a',20,'q'),('a',30,'r')) AS t(k,ord,v);
> {code}
> {noformat}
> ord COLLECT ARRAY_AGG COUNT
> 10 {p=1} [p] 1
> 20 {q=1} [p, q] 2
> 30 {r=1} [p, q, r] 3
> {noformat}
> {{ARRAY_AGG}} and {{COUNT}} accumulate over the window. {{COLLECT}} does not,
> in the same
> query on the same rows.
> The input here is already in ascending order, so this is not the out-of-order
> case -- it is
> wrong on ordinary input, and it fails quietly rather than throwing.
> Three comparisons on the same data, all of which do accumulate:
> {noformat}
> same query in BATCH mode {p=1} {p=1, q=1} {p=1, q=1, r=1}
> same query over a PROCTIME OVER window cumulative
> Apache Spark 4.2.0, collect_list ['p'] ['p','q'] ['p','q','r']
> {noformat}
> So it appears specific to the OVER window ordered by a non-time attribute.
> Reproduced on 2.1.3, 2.2.1 and 2.3.0.
> Reproducer attached.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)