Skip to content

fix: preserve native sliding integer sum overflow semantics - #6069

Open
rich7420 wants to merge 7 commits into
apache:mainfrom
rich7420:fix/6043-sliding-integer-sum
Open

rich7420 wants to merge 7 commits into
apache:mainfrom
rich7420:fix/6043-sliding-integer-sum

Conversation

@rich7420

@rich7420 rich7420 commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6043.

Rationale for this change

Native sliding integer sums currently wrap on overflow. Spark's ANSI sum throws and try_sum returns NULL, with recovery once the overflowing values leave the frame. This change preserves native window execution while matching those semantics.

What changes are included in this PR?

Add a retractable integer SUM accumulator for ANSI/TRY sliding windows. Two monotonic queues track the smallest and largest exact prefix sums, detecting intermediate overflow without rescanning each frame. Evaluate overflow after update and retraction so temporary overlap between frames cannot cause a false error. Use exact positive/negative retained-value bounds to avoid allocating prefix queues when overflow is impossible. Queue capacity is reported by the accumulator, but DataFusion 55.1 window operators do not account for accumulator sizes. Worst-case queue space remains linear in the largest non-null frame (including update-before-retract overlap), at 32 bytes per entry on a 64-bit host.

Disable reversal for ANSI/TRY sums because checked addition depends on input order. Legacy sliding sums retain the existing DataFusion path. Grouped and ever-expanding sums retain their existing accumulators.

Replace the proposed Spark fallback with native execution assertions. Extend SQL and Scala coverage for overflow, cancellation, recovery, NULL/empty frames, batch boundaries and RANGE peers. Extend the benchmark with DataFusion legacy SUM, Comet ANSI SUM and TRY_SUM controls, positive-only and near-limit inputs, and suffix frames. A bit_xor/count sink consumes the results without ANSI-dependent sink overflow. Since Spark recomputes suffix frames quadratically, verify those against Spark on up to 1,024 rows and compare the native modes on the full partition.

How are these changes tested?

Current upstream head is fbd8f7221831a118f55506e9655f1dab6a4876a2. Current-head upstream CI passed, including Required Checks. This commit adds the bound-based queue optimization, its native regressions and the expanded baseline benchmark.

The benchmark cases are implemented, but no new timing results for this exact head are attached here. The earlier Spark-fallback speedups below are not measurements of the latest optimization against DataFusion's legacy accumulator. That comparison remains to be reported, including the positive-only suffix workload and worst-case memory behavior.

After the upstream fix merges, prepare and independently validate the branch-1.1 backport requested in review.

Earlier verification is retained with its version:

  • At local candidate 3e8625af8: 70 Spark 4.1 Window/SQL tests, four focused Spark 3.5 tests with strict warnings, and ten Rust SUM tests passed. Rust coverage includes all 7,776 five-value sequences across four frame widths and three modes. Release build, workspace/all-target Clippy, formatting and license checks passed.
  • The twelve-case local benchmark used M3 Pro, Spark 4.1.3/JDK21, local[1], release native code, 65,536 rows, frame widths 16/1,024 and 0/50/100% NULL. It consumed SUM/TRY_SUM outputs and checked both execution paths. Native mean times were 27–33 ms versus 36–990 ms for Spark-window fallback, or 1.29–31.94x within that matrix. These are measurements of 3e8625af8, not fresh timings of the current head or production guarantees.
  • Fork #43 CI passed all Comet suites, Rust, strict compilation, benchmark checks, TPC-H/TPC-DS and all nine Spark 4.1 SQL shards at the earlier snapshot ce6e3a90d. It does not validate the new main merge.

The cancelled execution job in upstream run 36725911344 belongs to ce6e3a90d. It is historical evidence of an incomplete old run, not the current head's verdict; the successful current-head CI supersedes the old rerun request. Other suites skipped by a run are not covered by that run's verdict.

@github-actions github-actions Bot added the bug Something isn't working label Sep 21, 2026
@rich7420
rich7420 marked this pull request as ready for review September 21, 2026 07:09

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Correctness

Reviewed 46ef424d0261f73ea943af7160ca410c77c9854b against 6065705c16340c0be293212a71decfd9df4daae4. No verified P1/P2 findings.

The prior sliding path uses DataFusion's integer sum accumulator, which wraps on overflow without honoring the aggregate's ANSI/TRY mode. For a one-row-preceding frame over [Long.MaxValue, 1, -1], Spark's try_sum returns [Long.MaxValue, NULL, 0]; the wrapping path produces Long.MinValue at the second row. Spark recomputes changing frames, so the result can recover after the overflowing values leave.

The new admission guard correctly combines a lower bound other than UNBOUNDED PRECEDING, a LongType sum result, and a non-legacy aggregate evaluation mode. Spark promotes byte, short, int and long sums to long. Reading the mode through the existing shim also catches try_sum when session ANSI is disabled. Returning None follows the existing whole-WindowExec fallback contract, including shrinking, singleton and empty frames. Legacy integral sums, floating-point sums, and expanding integral sums retain their existing native paths; expanding sums use Comet's mode-aware accumulator.

Validation

  • Compared Spark's sum state and moving/growing/shrinking frame evaluation on the maintained 3.5 and 4.0 branches. Checked all four Comet version shims and the checksum-verified DataFusion 55.1.0 sliding accumulator.
  • The current expressions CI job actually ran windows/sliding_integer_sum.sql. Its nine fallback queries check the expected reason and Spark-equivalent answers; five ordinary queries assert native coverage; three error queries compare overflow failures. Coverage includes both overflow directions, recovery, null/empty frames, ROWS/RANGE, narrow integral inputs and native controls.
  • The current execution CI job passed the affected window tests. Both jobs checked out merge 908f9931454e, whose tree equals the reviewed head, and consumed the native artifact built by the same run with matching upload/download SHA-256 digests.
  • Local validation comprised source assertions and an independent signed-64 arithmetic model. The required fallback reason is absent from the base implementation, establishing the fixture's pre-fix failure at the source-contract level; I did not run the pre-fix Spark/JNI implementation. Current runtime evidence is Spark 4.1 CI. Canonical maintained 3.4 and 4.1 source branches were unavailable, so source-semantic coverage is limited to 3.5/4.0. The expressions job's one canceled test is an unrelated Spark-4.0-only regexp reproducer.

Performance

The extra work is a small planning-time type/mode check, with no per-row allocation or arithmetic added. Affected windows now execute in Spark, and one rejected expression causes the whole window operator to fall back. That can increase query time and JVM buffering/recomputation costs, including the quadratic shrinking-frame case, but it is the explicit cost of preventing incorrect overflow results. The guard is independent of input values and therefore also applies when a particular dataset would not overflow. Native legacy and expanding controls keep the unaffected cases covered. No performance measurement or speedup claim is attached to this change.

Design

The fix fits the existing sliding-decimal fallback and uses the same frame classification as native accumulator selection. Keeping the decision in window admission avoids changing grouped sums or pretending that a sticky overflow flag can support retraction. A native alternative would need to reproduce Spark's overflow and recovery semantics as frame contents change, which is a larger change than this compatibility guard. The fallback reason and compatibility documentation explain the resulting boundary.

Abstraction & complexity

The implementation reuses the existing evaluation-mode shim and fallback machinery. It introduces no new runtime abstraction or ownership path. The small test helper adapts existing plain-sum assertions to the session's ANSI mode, while the SQL fixture separately tests explicit TRY mode with ANSI both disabled and enabled. The focused production change, regression fixture and documentation form a coherent scope; I found no actionable simplification.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for this. I traced the premise through the native side and it checks out. process_agg_func in native/core/src/execution/planner.rs only builds the mode-aware SumInteger UDAF when is_ever_expanding is true, and sliding frames fall through to DataFusion's built-in sum, which wraps and has no notion of eval mode. So the divergence was real, and the admission guard is at the right layer, matching the decimal guard from #4732.

A few things I checked that look fine, noted here so they do not get re-litigated. s.dataType == LongType is the right predicate, because Spark promotes byte, short, int and long sums to LongType, and interval sums are already rejected earlier by AggSerde.sumDataTypeSupported. Shim coverage is not a concern either, since sumEvalMode already exists across the 3.4, 3.5, 4.0 and 4.1+ shims. I also suspected a matching gap in integral avg, but Spark's Average accumulates integral input in DoubleType, so there is no ANSI overflow there and nothing to fix.

The fixture is thorough on the behaviour it covers, including both overflow directions, recovery, null and empty frames, and native controls on the paths that stay native. Both of my comments are about the test harness rather than the fix itself.


-- Spark needs constant folding for PRECEDING bounds. Aggregate inputs remain columns.
statement
SET spark.sql.optimizer.excludedRules=

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the comment explaining why this is here, it made the intent easy to follow. One concern though. CometSqlFileTestSuite excludes ConstantFolding on purpose, and appends that exclusion after each file's own configs so that a header cannot override it. Clearing it with a SET statement turns folding back on for every query in the rest of the file, including the plain query entries that assert native coverage. This is the only fixture in the repo that does this.

window_functions.sql lines 198 to 203 hit the same ROWS ... N PRECEDING problem and took the other route, keeping SQL fixtures to bounds that parse directly and covering N PRECEDING in CometWindowExecSuite through the DataFrame API. Could we follow that convention here? The RANGE ... 1 PRECEDING cases already work unfolded, as window_functions.sql line 226 shows, so only the ROWS ones would need to move and the SET could go away entirely.

If you think the fixture really does need folding, could we add a supported file-level directive for it instead, so the harness stays in control of what it guarantees?

Seq("true", "false").foreach(aqeEnabled =>
withSQLConf(
// This native coverage matrix includes legacy sliding integral sums.
SQLConf.ANSI_ENABLED.key -> "false",

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I follow why this is needed. The default profile is Spark 4.1.3 where ANSI is on, so without this the sliding sums in this matrix would now fall back and checkSparkAnswerAndOperator would fail its native coverage assertion.

The side effect is that this matrix stops exercising ANSI at all, and it covers a good deal more than sums. Would routing the sum assertions through checkSlidingIntegralSum work here, the way the three tests below do, so the matrix keeps running in both modes? If pinning really is the cleaner option, could the comment mention that the matrix no longer covers ANSI? As written it reads as being only about legacy sums, and the next person may not realise the broader coverage went with it.


-- Config: spark.sql.adaptive.enabled=false
-- Config: spark.sql.ansi.enabled=false
-- Config: spark.comet.operator.WindowExec.allowIncompatible=false

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is spark.comet.operator.WindowExec.allowIncompatible=false doing anything here? CometWindowExec does not declare a support level, and the only other fixture using that key, lag_lead.sql, sets it to true. If this is just documenting the default, it may be worth dropping so that it does not get copied into future fixtures as boilerplate.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rechecked the unchanged 46ef424d after the new harness feedback. I agree with the existing comments: the fixture-wide SET clears the harness’s folding exclusion, the window matrix now forces legacy mode for non-sum cases too, and WindowExec.allowIncompatible has no effect on this conversion guard. Those threads cover the needed test changes. I found no additional P1/P2 to raise.

One coverage qualification: the matrix previously inherited the session’s ANSI mode. It was not an explicit two-mode matrix. The new fixture still tests ANSI/TRY fallback and errors, and its table-column sums are not folded away. Its five ordinary queries retain native-plan assertions.

Fresh log checks confirm that the SQL fixture and window suite passed on CI merge 908f9931, whose parents are the reviewed base/head and whose tree equals this head. This follow-up used source checks, with no new local Spark/JNI run. Maintained Spark 3.5/4.0 semantics were rechecked; the 3.4/4.1 source gaps remain. My existing approval is preserved.

@rich7420

Copy link
Copy Markdown
Contributor Author

cc @andygrove , @sunchao could you take another look? thanks!

@rich7420
rich7420 force-pushed the fix/6043-sliding-integer-sum branch from d7bd90a to 748bf99 Compare September 28, 2026 14:50

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@sunchao

sunchao commented Sep 28, 2026

Copy link
Copy Markdown
Member

@rich7420 actually instead of fallback to Spark which would sacrifice performance, should we consider to fix this behavior in Comet native path instead?

@rich7420
rich7420 marked this pull request as draft September 30, 2026 13:17
@github-actions github-actions Bot added the area:aggregation Hash aggregates, aggregate expressions label Sep 30, 2026
@rich7420 rich7420 changed the title fix: fall back for ANSI and TRY integer sums over sliding windows fix: preserve native sliding integer sum overflow semantics Sep 30, 2026
@rich7420

rich7420 commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor Author

Historical progress update for ce6e3a90d. Upstream #6069 has advanced to 67889aba5, whose CI passed. The old cancelled-run rerun request below is superseded; see the current PR description for head-specific validation.

Updated this PR with the native sliding accumulator from 3e8625af8. Current head ce6e3a90d merges main without changes to the SUM implementation or its regressions and benchmark. It checks ordered prefix extrema after retraction, so it detects intermediate overflow and recovers when the offending values leave the frame. ANSI/TRY sums cannot be reversed because that changes overflow behavior.

Fresh fork CI now passes at ce6e3a90d, exactly the current upstream head: Required Checks, all Comet suites including execution, Rust tests, strict Spark 3.5/JDK17 compilation, benchmark checks, TPC-H/TPC-DS and all nine Spark 4.1 SQL shards. This replaces the earlier fork verdict at 3e8625af8 with validation of the updated head.

Local Spark 4.1 Window/SQL tests (70), focused Spark 3.5 tests with strict warnings (4) and Rust SUM tests (10) passed at 3e8625af8. The 12-case local benchmark was faster than Spark-window fallback in every case, at 1.29–31.94x. The implementation, regression and benchmark diff is unchanged; the PR description records the local workload and limits.

This PR remains Ready for fresh review of the native implementation. Upstream run 36725911344 passed Rust, builds, strict compilation, scans/shuffle/expressions and TPC-H/TPC-DS. Its execution job was cancelled during CometJoinSuite, after the affected sliding integral ROWS tests passed in both ANSI modes. No ScalaTest failure was recorded before cancellation, but its execution-suite verdict remains incomplete and upstream Required Checks is still red. The successful fork run does not clear that required upstream check.

Could a maintainer rerun the cancelled upstream execution job and its dependent checks? This account cannot rerun the upstream workflow.

@rich7420
rich7420 marked this pull request as ready for review September 30, 2026 14:47
@andygrove

Copy link
Copy Markdown
Member

I would like to do a deeper review of this PR before it is merged.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I compared this against branch-1.1, where ANSI and TRY sliding integer sums run natively through DataFusion's wrapping sum, as well as against main. Moving from the fallback to a native retractable accumulator works out nicely: tracking the prefix-sum extrema reproduces Spark's recompute-every-frame semantics, including an intermediate overflow that later cancels and the recovery once those values leave the frame, and deferring the check to evaluate copes with DataFusion updating before it retracts. My earlier comments on the fixture and the window matrix are all addressed, so what follows is about measuring the new path against the one 1.1 users run today, plus a backport question.

1. spark/src/test/scala/org/apache/spark/sql/benchmark/CometSlidingSumBenchmark.scala:56

On Spark 4, ANSI is on by default, so after this change every sliding SUM over an integral column goes through SlidingSumIntegerAccumulator instead of DataFusion's SlidingSumAccumulator. That DataFusion path is what 1.1 runs today in every eval mode, so it is the baseline users upgrade from, but this benchmark only compares the new path with the Spark fallback. Could you add the DataFusion path as a case? On this branch a plain sum with spark.sql.ansi.enabled=false still uses DataFusion's accumulator, so running the same query in legacy and ANSI mode on one build measures the difference directly.

The data shape matters here too. The id % 17 - 8 values cancel out every 17 rows, which keeps the new queues short. All-positive values, such as amounts or quantities, make the prefix sums monotonic, and then one queue holds every non-null row in the frame. A suffix frame shows the worst case, because its first frame is the whole partition:

SELECT sum(s), count(s) FROM (
SELECT sum(v) OVER (ORDER BY id ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING) AS s
FROM (SELECT id, id % 17 AS v FROM range(4000000)))

Could you add positive-only values and a suffix frame like this one, with ANSI on and off?

2. native/spark-expr/src/agg_funcs/sum_int.rs:132

DataFusion's sliding accumulator keeps a running sum and a count. Here, when the prefix sums are strictly monotonic, as they are for all-positive data, minima or maxima keeps a 32-byte (usize, i128) entry for every non-null row in the frame, and VecDeque growth can double the capacity on top of that. For a suffix frame such as ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING, that is the whole partition. WindowAggExec and BoundedWindowAggExec in DataFusion 55.1 never call Accumulator::size(), so this memory isn't counted anywhere, even though size() reports it. Could the doc comment state that worst-case bound?

If the benchmark shows a cost on long frames, there is a way to keep ordinary data at constant memory. A partial sum that ends at a new row can't be larger than the sum of the positive values in the frame, or smaller than the sum of the negative values. If you track those two sums alongside start and end, an entry only needs to go into maxima once the positive sum exceeds i64::MAX, and into minima once the negative sum drops below i64::MIN. A skipped entry can never be the one that overflows in a later frame, and popping the entries it dominates stays safe, so both queues stay empty unless the frame holds values near the i64 limits.

3. General

1.1.0 runs these sliding sums natively through DataFusion's wrapping sum. So when a frame overflows, a Spark 4 user gets a wrapped value where Spark throws ARITHMETIC_OVERFLOW, and try_sum returns a wrapped value instead of NULL on every Spark version. Since that is a wrong answer in a shipped release, I think this fix should also go to branch-1.1. The cherry-pick looks clean to me. sum_int.rs and CometWindowExecSuite.scala are identical on branch-1.1 and on this PR's merge base, and the process_agg_func hunk applies to the same code there. Would you be up for opening the backport once this merges?

@rich7420

rich7420 commented Oct 7, 2026

Copy link
Copy Markdown
Contributor Author

Ran CometSlidingSumBenchmark at fbd8f72 with a release build (target-cpu=native), Spark 3.5.9, JDK 17 and an 8 GB heap on an Apple M3 Pro.

For 4M-row suffix frames, these are the ranges of mean timings across two independent JVM runs:

Input DataFusion legacy SUM Comet ANSI SUM Comet TRY_SUM
Positive, no NULLs 709–736 ms 512–546 ms 522–529 ms
Positive, 50% NULLs 790–811 ms 569–584 ms 568–593 ms
Near-limit, no NULLs 744–749 ms 571–581 ms 565–580 ms
All NULLs 535–585 ms 543–575 ms 546–578 ms

The positive/no-NULL case was about 1.35–1.39× faster than legacy SUM. All-NULL results were similar. The additional 65,536-row bounded-frame run showed smaller, noisier differences between native implementations.

All 112 timed cases completed, with the existing plan and bit_xor/count result assertions passing. Suffix validation compares 1,024 rows against Spark and all 4M rows between native implementations. These are end-to-end timings, not an isolated before/after measurement of the queue optimization; peak memory was not measured.

All runs exited successfully; Maven exec:java emitted a Hadoop shutdown classloader warning after completion.

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

Labels

area:aggregation Hash aggregates, aggregate expressions bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Sliding-window SUM(BIGINT) ignores ANSI and TRY overflow semantics

3 participants