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)

Reply via email to