xumingming opened a new pull request, #58926:
URL: https://github.com/apache/spark/pull/58926

   ### What changes were proposed in this pull request?
   
   Add a new logical optimizer rule, `RewriteSlidingFramesAsPrefixDifferences`, 
that rewrites a
   group of sliding `ROWS BETWEEN n PRECEDING AND CURRENT ROW` frames of `sum` 
sharing one
   `(PARTITION BY, ORDER BY)` key into a single running sum plus lag 
differences:
   
   ```sql
   -- Running total, reused as the base for every sliding frame
   C = sum(x) OVER (PARTITION BY k ORDER BY d
                   ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)
   
   -- An n-row sliding sum is the running total minus the running total n+1 
rows back
   sum(x) OVER (PARTITION BY k ORDER BY d
               ROWS BETWEEN n PRECEDING AND CURRENT ROW)
       = C - coalesce(lag(C, n + 1) OVER (PARTITION BY k ORDER BY d), 0)
   ```
   
   For a query with two rolling horizons over a serving log:
   
   ```sql
   SELECT
     sum(tokens) OVER (PARTITION BY user_id ORDER BY create_ts
       ROWS BETWEEN 6 PRECEDING AND CURRENT ROW)  AS tokens_7d,
     sum(tokens) OVER (PARTITION BY user_id ORDER BY create_ts
       ROWS BETWEEN 29 PRECEDING AND CURRENT ROW) AS tokens_30d
   FROM serving_log
   ```
   
   the logical plan goes from
   
   ```
   Window [sum(tokens) ... ROWS BETWEEN 6 PRECEDING AND CURRENT ROW AS 
tokens_7d,
           sum(tokens) ... ROWS BETWEEN 29 PRECEDING AND CURRENT ROW AS 
tokens_30d]
   ```
   
   to
   
   ```
   Window [(window_prefix_sum_0 - coalesce(lag(window_prefix_sum_0, -7, null)
              OVER (PARTITION BY user_id ORDER BY create_ts), 0)) AS tokens_7d,
           (window_prefix_sum_0 - coalesce(lag(window_prefix_sum_0, -30, null)
              OVER (PARTITION BY user_id ORDER BY create_ts), 0)) AS tokens_30d]
   +- Window [sum(tokens) OVER (PARTITION BY user_id ORDER BY create_ts
                ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS 
window_prefix_sum_0]
      +- ...
   ```
   
   (`-7` / `-30` is how `Lag` prints its frame offset, i.e. 7 / 30 rows back.)
   
   Properties that matter for review:
   
   * No new physical operator. The rule is logical only and the result is 
executed by the
     existing `WindowExec`, so there is nothing new to maintain in the 
execution layer.
   * The plan shape is preserved. Both window operators keep the same partition 
and order spec,
     so they share one sort and no exchange or sort is added (asserted by a 
test).
   * One running sum per distinct measure, plus a running count for nullable 
measures; several
     measures produce several groups, named `window_prefix_sum_<i>` / 
`window_prefix_count_<i>`.
   * Partial rewrites: members of the same window operator that do not qualify 
are left in place,
     not dropped.
   * A nullable measure gets a `count` guard (`CASE WHEN c - lag(c, w) = 0 THEN 
NULL ELSE ...`),
     so an all-NULL frame keeps returning NULL rather than 0.
   * Groups whose total frame width is below an internal gate 
(`MIN_TOTAL_FRAME_WIDTH`, 16 rows)
     are not rewritten; bench section D locates the naive-vs-rewrite crossover.
   * The rewrite is idempotent, and it runs before `CollapseWindow` in the 
`Operator combine`
     batch. It is also registered in `RuleIdCollection`.
   
   Safety is a default-deny certificate per aggregate class rather than an 
expression inspection:
   only `sum` over an integral (byte/short/int/long) measure in a legacy, 
ANSI-off session
   qualifies, because its accumulator is a `LongType` buffer updated with 
wrapping `Add`, i.e.
   arithmetic in Z/2^64Z where modular reduction commutes with subtraction. The 
difference is
   therefore bit-identical to the frame fold even when the running prefix, 
which can reach
   magnitudes no bounded frame ever reaches, overflows (covered by a test). 
Excluded: floating
   point (addition is not associative), decimal (an overflowed prefix silently 
becomes NULL under
   `CheckOverflowInSum`), anything running under ANSI (the newly reachable 
prefix magnitudes would
   throw), `try_sum`, distinct and FILTERed sums, non-`sum` aggregates, RANGE 
frames, growing
   frames, frames without ORDER BY, and window expressions whose spec differs 
from the window
   operator's own spec.
   
   The rewrite is opt-in and off by default:
   
   ```
   spark.sql.optimizer.windowPrefixRewrite.enabled = false   (new in 4.4.0)
   ```
   
   It declares `ConfigBindingPolicy.NOT_APPLICABLE`, since it changes only plan 
shape and does not
   affect view/UDF/procedure resolution.
   
   ### Why are the changes needed?
   
   Spark evaluates a bounded moving frame by re-aggregating every row inside it 
on every output
   row: `SlidingWindowFunctionFrame.write` resets the aggregate buffer and 
re-folds the whole
   frame. For nested sliding frames over one `(PARTITION BY, ORDER BY)` key the 
cost is the sum of
   the frame widths, 3+7+30+90+180 = 310 aggregate updates per output row per 
measure with the
   usual 3d/7d/30d/90d/180d horizon set, and it grows with every horizon added. 
There is no
   incremental path today: the window aggregate processor supports only
   `initialize`/`update`/`evaluate`, with no `remove`, so a moving frame cannot 
be advanced
   row-by-row. Growing frames are already incremental and offset frames are 
O(1); the sliding frame
   is the one remaining unbounded-per-row path. This shape is common in 
user-facing metric tables
   that compute several rolling horizons per measure, and in the workload that 
motivated this
   change the window operator was ~45% of the CPU of a 3,437 CPU-hour 
application, with the window
   stage's own input sort spilling ~14.8 TiB on the largest query.
   
   The rewrite makes every frame O(1) while leaving the plan's parallelism and 
sort structure
   untouched (`WindowPrefixRewriteBenchmark`, results committed alongside):
   
   | shape | sliding frames | segment tree (SPARK-56546) | prefix differences |
   |---|---|---|---|
   | A: 6 measures x widths 2/6/29/89/179 | 19,693 ms | 13,158 ms (1.5X) | 
**1,693 ms (11.6X)** |
   | B: 1 measure, total width 310 | 4,975 ms | 5,202 ms (1.0X) | **740 ms 
(6.7X)** |
   | C: 32,768 rows per partition | 16,492 ms | 10,105 ms (1.6X) | **2,459 ms 
(6.7X)** |
   
   This is complementary to the executor-side block-chunked segment tree 
(SPARK-56546) rather than
   a replacement: the segment tree is O(log W) per frame and accepts any 
`DeclarativeAggregate`
   without FILTER/DISTINCT, while this rewrite is O(1) per frame but restricted 
to invertible
   `sum` over integral, ANSI-off measures. With both enabled, a rewritten group 
has no moving
   frames left for the segment tree to handle.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes, a new feature, but no change for users who do not opt in.
   `spark.sql.optimizer.windowPrefixRewrite.enabled` is new and defaults to 
`false` (new in
   4.4.0, so there is no behavior change relative to any released version). 
When enabled,
   qualifying queries return exactly the same rows, the rewrite is 
bit-identical for the aggregate
   classes it accepts, but `EXPLAIN` shows the running sum plus lag-difference 
shape instead of the
   sliding frame, and the window operator performs O(1) work per frame instead 
of O(W).
   
   ### How was this patch tested?
   
   Two suites plus a benchmark, all added by this PR.
   
   `RewriteSlidingFramesAsPrefixDifferencesSuite` (catalyst, `PlanTest`) pins 
the emitted plans and
   every rejection path: the emitted running sum plus per-candidate lag offsets 
and group naming,
   several measure groups in one window operator, the NULL guard for nullable 
measures, the cost
   gate boundary (width 15 no, width 16 yes), partial rewrites (both for 
non-`sum` members and for
   `sum`s that fail the safety certificate), idempotence, that emitted frames 
are ROWS and never
   RANGE, and that no negative case is rewritten (filtered, distinct, float, 
decimal, `try_sum`,
   ANSI, min/max, no ORDER BY, RANGE, growing, spec mismatch).
   
   `WindowPrefixRewriteQuerySuite` (sql/core) differentially runs each query 
with the config off
   and on and compares the full projection, covering multiple 
widths/measures/nullability,
   duplicate `(PARTITION BY, ORDER BY)` keys, partitions shorter than the frame 
width, all-NULL
   measures, tinyint/smallint/int measures, 64-bit prefix overflow 
(bit-identical results), ANSI
   sessions, and that the rewrite adds no exchange or sort; it ends with a 
randomized differential
   test.
   
   ```
   build/sbt "catalyst/testOnly *RewriteSlidingFramesAsPrefixDifferencesSuite"
   build/sbt "sql/testOnly *WindowPrefixRewriteQuerySuite"
   build/sbt "hive/testOnly *SparkConfigBindingPolicySuite"
   ```
   
   The benchmark is runnable from a single command, and the committed results 
file is generated
   from it:
   
   ```
   SPARK_GENERATE_BENCHMARK_FILES=1 build/sbt "sql/test:runMain \
     org.apache.spark.sql.execution.benchmark.WindowPrefixRewriteBenchmark"
   ```
   
   ### Was this patch authored or co-authored using generative AI tooling?
   
   Generated-by: Pi
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to