[
https://issues.apache.org/jira/browse/FLINK-40737?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18123934#comment-18123934
]
sepuri sai krishna commented on FLINK-40737:
--------------------------------------------
While working on the fix I found the same root cause also breaks DISTINCT over
this kind of window.
COUNT(DISTINCT v) OVER (PARTITION BY k ORDER BY ord) returns wrong counts on
both state backends
once a value repeats. For rows (30,p),(10,r),(20,q),(25,r) it returns 1, 2, 3,
4 where the correct
counts are 1, 2, 2, 3.
The accumulator data view is bound to the partition key rather than to the sort
key, so every sort
key shares one view. That is the same reason COLLECT only ever shows the
current row, so one change
covers both and I have kept it in the same PR rather than splitting it.
> 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
> 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)