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?
Describe the bug
DataFusion's window operators register no
MemoryConsumer(nothing undersrc/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 atUNBOUNDED FOLLOWING, which includes Spark's default frame foragg(...) OVER (PARTITION BY k).WindowAggExeckeeps 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.BoundedWindowAggExecbuffers only one window partition at a time, but it doesn't reserve either.Spark's
WindowExecbuffers throughExternalAppendOnlyUnsafeRowArray, 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) ofPERCENT_RANK() OVER (PARTITION BY k ORDER BY v)over a skewed key should shownative_allocatedrunning 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
WindowAggExecis apache/datafusion#22946. Until that lands, could we reserve the batchesWindowAggExecbuffers, so an oversized partition fails the task instead of the executor?