Skip to content

[FLINK-40483][table-planner] Convert batch ROW_NUMBER Top-N to two-stage Rank - #29051

Draft
Timm0 wants to merge 1 commit into
apache:masterfrom
Timm0:FLINK-40483
Draft

[FLINK-40483][table-planner] Convert batch ROW_NUMBER Top-N to two-stage Rank#29051
Timm0 wants to merge 1 commit into
apache:masterfrom
Timm0:FLINK-40483

Conversation

@Timm0

@Timm0 Timm0 commented Aug 31, 2026

Copy link
Copy Markdown
Contributor

What is the purpose of the change

The batch planner compiles ROW_NUMBER() OVER (PARTITION BY … ORDER BY …) … WHERE rn <= N to a single-stage OverAggregate that hash-shuffles and sorts the entire input before dropping all but the top rows. This change routes that pattern through the existing two-stage Rank (a local top-N before the shuffle, and a global top-N after), so only the local survivors cross the Exchange. It is a rule-based conversion at the logical phase, aligning batch with streaming.

Brief change log

  • Relax the RANK-only guards in FlinkLogicalRankRuleForConstantRange and BatchPhysicalRankRule to also admit ROW_NUMBER, routing batch ROW_NUMBER() … WHERE rn <= N (constant range) to the two-stage Rank instead of OverAggregate
  • Add a RankType parameter to RankOperator, emit on the row-number counter for ROW_NUMBER, and on the rank counter for RANK
  • Add RankOperatorTest and RowNumberITCase regenerate the flipped ROW_NUMBER goldens in RankTest.xml and FlinkLogicalRankRuleForConstantRangeTest.xml
  • Add the ROW_NUMBER_TOP_N compiled-plan restore program + JSON

Verifying this change

  • Added RankOperatorTest and RowNumberITCase regenerate the flipped ROW_NUMBER goldens in RankTest.xml and FlinkLogicalRankRuleForConstantRangeTest.xml
  • Added the ROW_NUMBER_TOP_N compiled-plan restore program + JSON

I also ran some benchmark tests on a dataset with near-unique keys (1 row per key) and lot's of duplicated keys (~1800 rows per key). Below are the results:

  1. ROW_NUMBER() OVER (PARTITION BY ... ORDER BY ... DESC) AS rn ... WHERE rn = 1

Near-unique Keys:

  • performance decrease of ~12%
  • exchanged rows decreased by ~10%

Duplicated Keys:

  • performance increase of ~48%
  • exchanged rows decreased by ~10%
  1. ROW_NUMBER() OVER (PARTITION BY ... ORDER BY ... DESC) AS rn ... WHERE rn <= 200

Near-unique Keys:

  • performance increase of ~7%
  • exchanged rows decreased by ~1%

Duplicated Keys:

  • performance decrease of ~2%
  • exchanged rows decreased by ~1%

Does this pull request potentially affect one of the following parts:

  • Dependencies (does it add or upgrade a dependency): no
  • The public API, i.e., is any changed class annotated with @Public(Evolving): no
  • The serializers: no
  • The runtime per-record code paths (performance sensitive): yes
  • Anything that affects deployment or recovery: JobManager (and its components), Checkpointing, Kubernetes/Yarn, ZooKeeper: no
  • The S3 file system connector: no

Documentation

  • Does this pull request introduce a new feature? no
  • If yes, how is the feature documented? na

Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Opus 4.8 (1M context)

…age Rank

- Relax the `RANK`-only guards in `FlinkLogicalRankRuleForConstantRange` and `BatchPhysicalRankRule` to also admit `ROW_NUMBER`, routing batch `ROW_NUMBER() … WHERE rn <= N` (constant range) to the two-stage `Rank` instead of `OverAggregate`
- Add a `RankType` parameter to `RankOperator`, emit on the row-number counter for `ROW_NUMBER`, and on the rank counter for `RANK`
- Add `RankOperatorTest` and `RowNumberITCase` regenerate the flipped `ROW_NUMBER` goldens in `RankTest.xml` and `FlinkLogicalRankRuleForConstantRangeTest.xml`
- Add the `ROW_NUMBER_TOP_N` compiled-plan restore program + JSON
@Timm0
Timm0 marked this pull request as ready for review August 31, 2026 15:09
@flinkbot

flinkbot commented Aug 31, 2026

Copy link
Copy Markdown
Collaborator

CI report:

Bot commands The @flinkbot bot supports the following commands:
  • @flinkbot run azure re-run the last Azure build

@Timm0
Timm0 marked this pull request as draft August 31, 2026 15:17
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants