Skip to content

Native window operators reserve no memory for the batches they buffer #6253

Description

@andygrove

Describe the bug

DataFusion's window operators register no MemoryConsumer (nothing under src/windows/ in datafusion-physical-plan 55.1.0 reserves memory). So whatever Comet's native windows buffer is invisible to Comet's memory pool and to Spark.

The worst case is WindowAggExec. The planner uses it whenever a window expression can't run in bounded memory (planner.rs#L2420-L2446): PERCENT_RANK, CUME_DIST, NTILE, and any aggregate frame that ends at UNBOUNDED FOLLOWING, which includes Spark's default frame for agg(...) OVER (PARTITION BY k). WindowAggExec keeps every input batch of the task's partition (window_agg_exec.rs:606) and concatenates them before evaluating (window_agg_exec.rs:539). A large or skewed partition can then hold a multiple of its input size natively with no reservation behind it. BoundedWindowAggExec buffers only one window partition at a time, but it doesn't reserve either.

Spark's WindowExec buffers through ExternalAppendOnlyUnsafeRowArray, which spills. The native path instead grows until the executor hits its container limit, and the pool never gets a chance to refuse anything. A string column that concatenates past 2 GiB fails the concat outright.

Steps to reproduce

I haven't measured this yet. A traced run (spark.comet.tracing.enabled=true) of PERCENT_RANK() OVER (PARTITION BY k ORDER BY v) over a skewed key should show native_allocated running ahead of the pools' reserved total by about the size of the largest partition.

Expected behavior

Native windows reserve what they buffer. A partition that doesn't fit then fails the task with a memory error, or spills once DataFusion supports it, instead of growing unaccounted.

Additional context

Spilling for WindowAggExec is apache/datafusion#22946. Until that lands, could we reserve the batches WindowAggExec buffers, so an oversized partition fails the task instead of the executor?

Activity

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

Metadata

Metadata

Assignees

Labels

area:memoryMemory pools, reservations, OOM handlingbugSomething isn't workingrequires-triage

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions