[ 
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)

Reply via email to