James Xu created SPARK-59658:
--------------------------------
Summary: Rewrite sliding window frames as prefix differences
Key: SPARK-59658
URL: https://issues.apache.org/jira/browse/SPARK-59658
Project: Spark
Issue Type: Improvement
Components: SQL
Affects Versions: 4.4.0
Reporter: James Xu
Spark executes a bounded moving window frame by re\-aggregating every row
currently inside the frame on every output row: {{SlidingWindowFunctionFrame}}
resets the aggregate buffer and re\-folds the whole frame for each row
produced. For a window operator with nested sliding frames over one \(PARTITION
BY, ORDER BY\) key, the cost is the sum of the frame widths \(Σw\)
aggregate\-update evaluations per output row, and the cost is structural —
every added horizon multiplies the work. There is no incremental path: the
window aggregate processor only supports
{{initialize}}/{{update}}/{{evaluate}}, with no {{remove}}, so a moving frame
cannot be advanced row\-by\-row.
The pattern is common in user\-facing metric tables that compute several
rolling horizons per measure over a daily partition ordering. A realistic
example:
{code:sql}
SELECT
seller\_user\_id,
sum\(cnt\) OVER \(PARTITION BY seller\_user\_id ORDER BY receive\_date
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW\) AS cnt\_3d,
sum\(cnt\) OVER \(PARTITION BY seller\_user\_id ORDER BY receive\_date
ROWS BETWEEN 6 PRECEDING AND CURRENT ROW\) AS cnt\_7d,
sum\(cnt\) OVER \(PARTITION BY seller\_user\_id ORDER BY receive\_date
ROWS BETWEEN 29 PRECEDING AND CURRENT ROW\) AS cnt\_30d,
sum\(cnt\) OVER \(PARTITION BY seller\_user\_id ORDER BY receive\_date
ROWS BETWEEN 89 PRECEDING AND CURRENT ROW\) AS cnt\_90d,
sum\(cnt\) OVER \(PARTITION BY seller\_user\_id ORDER BY receive\_date
ROWS BETWEEN 179 PRECEDING AND CURRENT ROW\) AS cnt\_180d
\-\- ...repeated for several more measures
FROM daily\_aggregates
{code}
This shape measures 5 horizons, i.e. Σw = 3\+7\+30\+90\+180 = 310
aggregate\-update evaluations per output row per measure family.
In a production workload measured against this shape \(a 3,437 CPU\-hour
application\), 45% of the CPU sits inside the window operator, and the window
stage's own input\-sort spill reaches 14.8 TiB on the largest query.
h3. Root cause
{{SlidingWindowFunctionFrame.write}} re\-runs the aggregate update for every
row in the frame on every output row. Growing frames are incremental and offset
frames are O\(1\), so the sliding frame is the only path with no O\(1\)
evaluation. The work factor is exactly Σw \(sum of frame widths\) per group per
row.
h3. Proposed change
Add an optimizer rule that rewrites each group of qualifying sliding frames as
one running aggregate plus offset differences:
{code}
C = sum\(x\) OVER \(PARTITION BY k ORDER BY d ROWS BETWEEN
UNBOUNDED PRECEDING AND CURRENT ROW\)
sum\(x\) OVER \(... ROWS BETWEEN n PRECEDING AND CURRENT ROW\)
= C \- coalesce\(lag\(C, n \+ 1\) OVER \(PARTITION BY k ORDER BY d\), 0\)
{code}
For integral accumulators with ANSI off this is provably bit\-identical: the
original fold already computes with wrapping addition, and modular reduction
commutes with subtraction, so the running prefix — a magnitude no bounded frame
can reach — overflows harmlessly. Floating\-point accumulators
\(non\-associative addition\) and decimal accumulators \(overflow becomes a
silent NULL\) are excluded, and a nullable measure carries a running count so
an all\-NULL frame still returns NULL. The rewrite reuses the window operator's
own partition and order spec, so no new exchange or sort is introduced.
The rule is conservative: only non\-distinct, unfiltered {{sum}} over integral
measures with literal {{n PRECEDING AND CURRENT ROW}} frames qualifies, groups
below a minimum total frame width are skipped, and it is gated behind an
opt\-in configuration flag, default off.
h3. Expected impact* Per\-group per\-row work drops from Σw to H \+ 2 aggregate
evaluations \(H = number of distinct widths\): 310 → 7 \(44x\) for the example
shape.
* Microbenchmark on the production shape \(6 measures, widths 2/6/29/89/179, 2M
rows\): 19.7s → 1.69s per window pass \(11.6x\); results verified
bit\-identical against the original plan, including under prefix overflow past
2\^63.
* On the real production query shape, a measured A/B of the rewritten form
shows \~10.5x end\-to\-end and \~20x on the window operator alone.
* Production model for one risk\-model workload family \(21 nodes, window stage
4,424 CPU\-hours\): \~3,481 CPU\-hours recoverable \(projected\).
* No new exchange or sort: both window operators share one sort, and
Exchange/SortExec counts are unchanged after the rewrite.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]