Skip to content

perf: Partitioned topk emit row number - #26133

Open
SubhamSinghal wants to merge 3 commits into
apache:mainfrom
SubhamSinghal:partitioned-topk-emit-row-number
Open

SubhamSinghal wants to merge 3 commits into
apache:mainfrom
SubhamSinghal:partitioned-topk-emit-row-number

Conversation

@SubhamSinghal

@SubhamSinghal SubhamSinghal commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Rationale for this change

With enable_window_topn, WindowTopN rewrites Filter(rn <= K) → BoundedWindowAggExec → SortExec into:

BoundedWindowAggExec: row_number() PARTITION BY [pk] ORDER BY [val], mode=Sorted
  PartitionedTopKExec: fn=row_number, fetch=2, partition=[pk], order=[val ASC]

The window node stays only to produce rn, but the operator already knows every retained row's rank when it emits it: rows leave in (partition keys, order keys) order, and the retained set is exactly the rows with rank ≤ K. The node's cost scales with the operator's output (K × partitions) and never amortizes, because every batch it sees is dense in partition boundaries.

What changes are included in this PR?

  • PartitionedTopKExec::try_new takes ranking_field: Option<FieldRef>. When set, the operator appends a UInt64 column with each row's rank, and WindowTopN drops the window node:

    PartitionedTopKExec: fn=row_number, fetch=2, partition=[pk], order=[val ASC], emit=[rn]
    
  • Each policy derives the rank from what it already holds. ROW_NUMBER: position in the partition's drain. RANK: position of the first row with the same key; boundary ties take the boundary key's. DENSE_RANK: index of the distinct-key group. All three rely on the retained set being a complete order-prefix of its partition.

  • Applied only when the window node has exactly one expression, since the operator appends one column. Otherwise the node is kept, as today. The output schema is the node's (input fields plus that field), so column indices above are unchanged and schema_check() holds.

  • compute_properties declares the ordering the removed node published, [partition keys..., rank ASC NULLS LAST], so ORDER BY pk, rn stays sort-free.

Are these changes tested?

  • topk/mod.rs: the rank column for each policy, including boundary ties, an eviction that moves a row to the tie list, a DENSE_RANK group gathered from several input batches, and output spanning several batch_size chunks.
  • sorts/partitioned_topk.rs: the widened schema, and the [pk, rn] ordering (but not [rn] alone).
  • core/tests/physical_optimizer/window_topn.rs: the node is removed for one expression and kept for two; ORDER BY pk, rn plans no SortExec for all three functions; end to end, the emitted column equals the flag-off result (ties, 4 partitions, batch size 4).
  • window_topn.slt and range_partitioning.slt: EXPLAIN output updated, query results unchanged.

Benchmark

h2o window.sql k=2 sweeps on J1_1e7_1e7_NA.parquet (10 M rows, 14 partitions), enable_window_topn = true in both arms.

partitions ROW_NUMBER RANK DENSE_RANK
100 0.037 → 0.035 s (1.06×) 0.036 → 0.037 s (0.97×) 0.037 → 0.037 s (1.00×)
1 K 0.036 → 0.035 s (1.03×) 0.038 → 0.039 s (0.97×) 0.040 → 0.040 s (1.00×)
10 K 0.041 → 0.038 s (1.08×) 0.034 → 0.032 s (1.06×) 0.047 → 0.044 s (1.07×)
100 K 0.075 → 0.050 s (1.50×) 0.092 → 0.060 s (1.53×) 0.147 → 0.118 s (1.25×)

Open question for reviewers

WindowTopN recognises ranking functions by name. Now that the window node is removed, a user-registered rank / row_number / dense_rank never runs its own evaluator:

  • if it returns a type other than UInt64, the query fails at execution (column types must match schema types);
  • if it returns UInt64, its values are silently replaced by the built-in's.

Which fix do you prefer?

  • A. Match by identity: == rank_udwf() etc. (WindowUDF equality compares the concrete type and its fields). This needs datafusion-functions-window as a normal dependency of datafusion-physical-optimizer, which adds no crates for anyone already using datafusion.
  • B. A WindowUDFImpl method declaring top-K ranking semantics, like limit_effect. It is a public API addition, so probably a follow-up.
  • C. Emit the column only when the field is UInt64. This fixes the error but not the silent replacement.

Are there any user-facing changes?

Only with enable_window_topn (default false): EXPLAIN no longer shows the window node for single-expression ranking queries, and PartitionedTopKExec shows emit=[...]. Query results are unchanged.

PartitionedTopKExec::try_new gains a parameter, so this is a breaking API change; please add the api change label. PartitionedTopKExec::ranking_field() is the new accessor.

@github-actions github-actions Bot added optimizer Optimizer rules core Core DataFusion crate sqllogictest SQL Logic Tests (.slt) physical-plan Changes to the physical-plan crate labels Oct 8, 2026
@codecov-commenter

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 87.88927% with 70 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.75%. Comparing base (1c43063) to head (318d470).
⚠️ Report is 14 commits behind head on main.

Files with missing lines Patch % Lines
datafusion/physical-plan/src/topk/mod.rs 89.84% 6 Missing and 33 partials ⚠️
...fusion/physical-plan/src/sorts/partitioned_topk.rs 83.97% 9 Missing and 20 partials ⚠️
datafusion/physical-optimizer/src/window_topn.rs 83.33% 0 Missing and 2 partials ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #26133      +/-   ##
==========================================
+ Coverage   82.73%   82.75%   +0.02%     
==========================================
  Files        1147     1147              
  Lines      449389   450243     +854     
  Branches   449389   450243     +854     
==========================================
+ Hits       371788   372594     +806     
+ Misses      54943    54932      -11     
- Partials    22658    22717      +59     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@SubhamSinghal
SubhamSinghal marked this pull request as ready for review October 9, 2026 04:44
@SubhamSinghal

Copy link
Copy Markdown
Contributor Author

@kosiew @jayzhan-synnada @kumarUjjawal Can you help in reviewing this PR

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

core Core DataFusion crate optimizer Optimizer rules physical-plan Changes to the physical-plan crate sqllogictest SQL Logic Tests (.slt)

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants