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)