Skip to content

Support WindowGroupLimitExec (window-based limit pushdown) #4837

Description

@andygrove

Spark inserts WindowGroupLimitExec to push a rank/row-number limit below a window (for example WHERE rn <= 10 over a ranking window). Comet does not support this operator today, so plans containing it fall back to Spark.

It is currently marked as planned in the operator support reference under Window.

Activity

  1. self-assigned this
    on Jul 8, 2026
  2. comphead commented on Jul 8, 2026

    @comphead
    Contributor

    related to #2551

  3. comphead commented on Jul 8, 2026

    @comphead
    Contributor

    This operator is a new TopK operator introduced in Spark 3.5

    Before Spark 3.5

     Catalyst rewrote this to something like Limit(3, Window(Sort(t))), which combined with TakeOrderedAndProjectExec gave you the classic heap-based top-K over a single stream. But this only worked when there
      was no PARTITION BY — the moment you added one, you fell back to the naive full-sort-full-window plan.
    
  4. comphead commented on Jul 8, 2026

    @comphead
    Contributor

    Prob we can use PartitionedTopKExec DF operator to handle this

  5. comphead commented on Jul 24, 2026

    @comphead
    Contributor

    we currently addressing rank in Comet but we can migrate to apache/datafusion#22885
    and dense_rank apache/datafusion#23869

  6. added this to the 1.1.0 milestone on Aug 1, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

Type

No type

Projects

No projects

    Milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions