Sumaiya Azad created FLINK-40964:
------------------------------------
Summary: Streaming Top-N over an updating ROW_NUMBER input never
retracts a replaced row (outer Rank uses AppendFastStrategy)
Key: FLINK-40964
URL: https://issues.apache.org/jira/browse/FLINK-40964
Project: Flink
Issue Type: Bug
Components: Table SQL / Planner
Affects Versions: 2.2.1
Environment: Flink 2.2.1 SQL client on a session cluster (Java 17);
PyFlink 2.2.1 (Java 21)
Reporter: Sumaiya Azad
Attachments: flink_topn_stale_row.py, flink_topn_stale_row.sql
h3. What happens
We keep the latest value per key, then pick the top key by that value. When a
key's latest value changes, streaming mode never retracts the key's old row, so
the result keeps the old value. Batch mode gives the right answer.
h3. How to reproduce
Four rows. Key A first has a value of 100, then a later value of 1. Key B has
50. The last row is a redelivered copy.
{code:sql}
-- t(id STRING, seq INT, k STRING, ts TIMESTAMP(3), v DOUBLE)
-- ('a1', 1, 'A', '2024-06-01 08:00:00', 100.0)
-- ('b1', 2, 'B', '2024-06-01 08:05:00', 50.0)
-- ('a2', 3, 'A', '2024-06-01 08:10:00', 1.0)
-- ('a2', 4, 'A', '2024-06-01 08:10:00', 1.0)
WITH latest AS ( -- latest value per key
SELECT k, v FROM (
SELECT *, ROW_NUMBER() OVER (PARTITION BY k ORDER BY ts DESC, id DESC) AS rk
FROM t) WHERE rk = 1)
SELECT k, v, r FROM ( -- top 1 key by that value
SELECT k, v, ROW_NUMBER() OVER (ORDER BY v DESC, k ASC) AS r FROM latest)
WHERE r <= 1
{code}
Two attached files run this query in streaming mode and in batch mode:
* {{{}flink_topn_stale_row.sql{}}}, for the Flink SQL client:
{{./bin/sql-client.sh -f flink_topn_stale_row.sql}}
* {{{}flink_topn_stale_row.py{}}}, for PyFlink: {{{}python
flink_topn_stale_row.py{}}}. It also prints the plan and exits 0 when the bug
reproduces.
Reproduced two ways on 2026-09-28: with the SQL client on a Flink 2.2.1 session
cluster (Java 17), and with PyFlink 2.2.1 (Java 21).
h3. Expected
{{{}('B', 50.0, 1){}}}. A's latest value is 1, so B is on top.
h3. Actual
* Streaming mode: the changelog is a single {{{}+I ('A', 100.0, 1){}}}. A's
old value of 100 is never removed.
* Batch mode: {{{}('B', 50.0, 1){}}}, which is correct.
h3. Likely cause
The streaming plan runs both ranks with {{{}AppendFastStrategy{}}}:
{code:java}
Rank(strategy=[AppendFastStrategy], ... orderBy=[v DESC, k ASC], select=[k, v])
+- Rank(strategy=[AppendFastStrategy], ... partitionBy=[k], orderBy=[$4 DESC,
id DESC])
{code}
{{AppendFastStrategy}} assumes its input only ever inserts rows. The inner rank
does not: it replaces A's row when a later row for A arrives. The outer rank,
therefore, never retracts A's old row.
h3. Why it matters
Any Top-N or filter built on a "latest row" query can publish a wrong row, and
nothing reports an error. On real data, at a parallelism of 2, we saw this in 8
such queries, even with events delivered in time order.
First reported on [email protected]:
[https://lists.apache.org/thread/1zw7fd76ypks7qgryn681qxj66z075ls]
--
This message was sent by Atlassian Jira
(v8.20.10#820010)