Skip to content

Support equality with one in DF55 WindowTopN safely - #127

Closed
osipovartem wants to merge 1 commit into
embucket-sync-df55.0.0from
df55-window-topn-eq-one
Closed

osipovartem wants to merge 1 commit into
embucket-sync-df55.0.0from
df55-window-topn-eq-one

Conversation

@osipovartem

Copy link
Copy Markdown
Collaborator

What changes

  • Recognize rank_col = 1 and 1 = rank_col in the opt-in DF55 WindowTopN rule; leave other equalities unchanged.
  • Reject rewrites when a sibling window expression (for example LEAD) depends on pruned rows.
  • Reject rewrites with no effective order key, including ORDER BY on the partition key itself. This also closes a pre-existing panic path for inequality predicates.
  • Add result and physical-plan regressions for equality, RANK ties, sibling LEAD, missing/redundant order keys, and unsupported literals.

Upstream counterpart: apache#26160. The DF55 branch needs the sibling and effective-order guards that already exist in current upstream.

Validation

  • Focused window_topn.slt SQLLogicTest passed.
  • cargo +1.95.0 test --profile ci -p datafusion-physical-optimizer --lib: 34 passed.
  • cargo +1.95.0 clippy --profile ci -p datafusion-physical-optimizer --all-targets -- -D warnings passed.
  • cargo +1.95.0 fmt --all --check passed.
  • Independent read-only review approved correctness, regressions, planner cost, API compatibility, and tests after two blocking findings were fixed.

Performance blocker

Five local 100k-row CLI runs, median query elapsed (opt-in rule off/on):

Partitions Off On
100 33 ms 27 ms
90,000 123 ms 644 ms

The high-cardinality case is about 5.2x slower. This PR remains draft and must not be merged or pinned into Rustice until that regression has an acceptable mitigation and benchmark evidence. The default enable_window_topn flag remains false. A pushed-down FilterExec projection can also block the rewrite; this change does not claim a Snowplow speedup.

@osipovartem

Copy link
Copy Markdown
Collaborator Author

Likely implementation-level source of the measured high-cardinality regression (not yet a CPU profile): DF55 PartitionedTopK::insert_batch builds a separate take_record_batch and heap entry for each partition seen in a batch, then size() sums every heap on each batch. At 90k near-unique groups, that creates many tiny batches/heaps and repeated accounting work. Current upstream instead gathers admitted rows once per input batch into a shared store and maintains O(1) running allocation totals.

The right follow-up is to benchmark and port the shared-store approach (or an equivalent bounded-memory design) with focused high-cardinality correctness, memory, and spill tests. Avoid enabling this draft rule in Rustice before that evidence exists. The preceding five-run 123/644 ms off/on result is the acceptance blocker.

@osipovartem

Copy link
Copy Markdown
Collaborator Author

Closing this draft. Rustice already registers its own RowNumberTopK rule before the DataFusion defaults and uses the spillable GroupedTopKExec for ROW_NUMBER() = 1. This generic DF55 port is not needed for the Snowplow path, and the separate datafusion-cli benchmark found a 5.2x high-cardinality slowdown when its opt-in rule fires. We should not merge or pin an extra competing path without a distinct use case and acceptable performance evidence. The generic upstream contribution remains open as apache#26160.

@osipovartem osipovartem closed this Oct 9, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant