Timo Theusner created FLINK-40483:
-------------------------------------

             Summary: Support batch Top-N optimization for ROW_NUMBER()
                 Key: FLINK-40483
                 URL: https://issues.apache.org/jira/browse/FLINK-40483
             Project: Flink
          Issue Type: Improvement
          Components: Table SQL / Planner
            Reporter: Timo Theusner


The batch planner does not optimize {{ROW_NUMBER() OVER (PARTITION BY … ORDER 
BY …) … WHERE rn <= N}} into a Top-N Rank operator. It falls back to a full 
{{OverAggregate}} that computes the row number for every row, followed by a 
Calc filter. The streaming planner does perform this optimization 
(Rank/Deduplicate).

Example query:
{code:SQL}
SELECT a, b, c FROM (
  SELECT a, b, c, ROW_NUMBER() OVER (PARTITION BY b ORDER BY c DESC) AS rn
  FROM MyTable) t
WHERE rn = 1
{code}
Batch plan:
{code:SQL}
Calc(where=[w0$o0 = 1])
+- OverAggregate(ROW_NUMBER(*) … UNBOUNDED PRECEDING .. CURRENT ROW)
   +- Sort(b ASC, c DESC)
      +- Exchange(hash[b])
{code}
Streaming plan:
{code:SQL}
Rank(rankType=[ROW_NUMBER], rankRange=[1,1], partitionBy=[b], orderBy=[c DESC])
{code}





--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to