Skip to content

fix: match Spark's ANSI integral SUM overflow error - #6764

Open
0lai0 wants to merge 1 commit into
apache:mainfrom
0lai0:fix-6661-sum-int-overflow-message
Open

0lai0 wants to merge 1 commit into
apache:mainfrom
0lai0:fix-6661-sum-int-overflow-message

Conversation

@0lai0

@0lai0 0lai0 commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6661

Rationale for this change

With spark.sql.ansi.enabled=true, Comet raised ARITHMETIC_OVERFLOW on an integral SUM overflow with the message integer overflow and an empty alternative. Spark reports long overflow with the suggestion Use 'try_add' to tolerate overflow and return NULL instead.

Spark's SUM over integral input always returns LONG and adds through Add(left, right, evalContext). The overflow therefore goes through MathUtils.addExact(Long, Long, context), which passes try_add as the hint, whatever the input type. Comet already throws in the same cases. Only the message parameters differ, and that matters to code that inspects getMessageParameters() or the message text.

What changes are included in this PR?

  • Add long_add_overflow_error() to native/spark-expr/src/lib.rs. It returns SparkError::ArithmeticOverflow { from_type: "long", function_name: "try_add" }, using the function_name field that fix: match Spark ANSI integral overflow errors (#6217) #6249 added.
  • Use it at the three overflow sites in sum_int.rs: the ungrouped update (the ungrouped merge reuses it), the grouped update, and the grouped merge. The other callers of arithmetic_overflow_error are unchanged.
  • planner.rs also uses SumInteger for window SUM over ever-expanding frames, so those windows now report the Spark error too. Sliding frames go through DataFusion's built-in sum and are not changed here. That path still ignores ANSI and returns a wrapped value instead of throwing, which Sliding-window SUM(BIGINT) ignores ANSI and TRY overflow semantics #6043 tracks.

How are these changes tested?

  • Unit tests in sum_int.rs pin from_type == "long" and function_name == "try_add" on every overflow path.
  • The two ANSI SUM tests in CometAggregateSuite used to pass with or without the fix. They now use checkSparkError and check that alternative carries try_add. They cover overflow and underflow, grouped and ungrouped queries, and overflow in the partial aggregate and in the final merge.
  • Both Scala tests pass under -Pspark-3.4 and -Pspark-4.1, and fail with the native change reverted.

@github-actions github-actions Bot added bug Something isn't working area:aggregation Hash aggregates, aggregate expressions area:expressions Expression evaluation labels Oct 7, 2026

@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.

No introduced P1/P2 issues found within this review.

Reviewed all three files in the full PR diff at d27cd8dad121caf25ace3b27b61f7ef443f07422, against supplied base b56349697b786ff2ad1c1bcf6ecf45b809af5f30 using merge-base ba7c89203813b848fb620f1078276b9f42614f96. GitHub confirms this scope and non-draft status. The snapshot and live discussion contain no existing reviews, conversation comments, inline comments, or threads.

Routed skills: review-comet-pr and review-comet-expression-pr. Read repository AGENTS.md and relevant contributor guidance.

Summary

  • Prior state and problem: ANSI integral SUM already detected overflow, but reported integer overflow without Spark’s try_add suggestion. This affected consumers inspecting exception messages and parameters.
  • Design approach: Add long_add_overflow_error() and use it at the three existing ANSI overflow sites. The helper supplies from_type = "long" and function_name = "try_add" through the existing structured-error pipeline.
  • Correctness: Traced ungrouped update and merge, grouped update, and grouped merge. Spark’s integral SUM widens to LongType and uses Add for update and merge. Its checked addition supplies try_add. All 13 focused Rust accumulator tests passed, including the six new error-parameter tests. No introduced P1/P2 correctness issue was identified.
  • Compatibility analysis: Compared Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0, including Sum, arithmetic evaluation/code generation, MathUtils, and error formatting. Spark 4.2 canonicalizes long overflow to overflow. Comet’s existing 4.x shim invokes that formatter, preserving the version-specific behavior. No schema, configuration, or support-level change is introduced.
  • Key design decisions: Keeping the new constructor separate leaves other callers of arithmetic_overflow_error unchanged. Legacy wrapping, TRY-mode null behavior, and decimal aggregation retain their existing implementations.
  • Implementation sketch: Replace three error constructors, add native assertions for overflow and underflow, and strengthen two Scala tests. The Scala setup controls one versus two input partitions to exercise partial aggregation and final merging, both grouped and ungrouped.
  • Performance: Successful aggregation retains the same arithmetic, loops, and accumulator state. The additional suggestion string is constructed only on failure. No benchmark was run, and no evidence-backed performance regression was identified.
  • Design: The change directly repairs error parity without changing aggregation routing or execution. Its scope is small and easy to trace.
  • Abstraction & complexity: One named constructor keeps the three sites consistent without introducing a new framework or expanding the generic helper’s responsibilities. No P1/P2 complexity concern was identified.
  • Behavioral changes worth calling out: Compared the affected path with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1. The error-detail correction is intentional. Ever-expanding window sums using SumInteger also inherit it. The separate sliding-window path is unchanged.
  • Suggested improvements: No additional code change met the P1/P2 reporting bar. CI and JVM integration execution remain outstanding validation, as detailed below.

Exact-head CI: Only the labeling check succeeded. Comet CI, CodeQL, and both title-check runs report action_required. Comet CI has zero jobs, so there is no build or test verdict.

Validation: cargo test --offline --locked -p datafusion-comet-spark-expr agg_funcs::sum_int::tests --lib passed all 13 tests. git diff --check passed and tracked files remain unchanged. I inspected but did not run the Scala integration tests or Spark SQL suites. This checkout lacks built JNI/JVM artifacts and a Spark Maven dependency cache, so runtime parity across JVM profiles remains unverified locally.

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 area:expressions Expression evaluation bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

ANSI integral SUM overflow reports "integer overflow" without Spark's try_add suggestion

2 participants