Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@

## sum (window)

- 4.1.1, audited 2026-06-25: `SUM(<decimal>)` over a sliding window frame (lower bound not `UNBOUNDED PRECEDING`) falls back to Spark. The native sliding path would use DataFusion's built-in `sum`, which wraps on overflow rather than returning Spark's NULL, and overflow cannot be detected at plan time, so the whole sliding decimal `SUM` case falls back. Ever-expanding frames use Comet's overflow-aware `SumDecimal` UDAF and run natively, matching Spark including overflow-to-NULL. Bigint sliding `SUM` overflow matches Spark (both wrap) and stays native.
- 4.1.1, audited 2026-06-25: `SUM(<decimal>)` over a sliding window frame (lower bound not `UNBOUNDED PRECEDING`) falls back to Spark. The native sliding path would use DataFusion's built-in `sum`, which wraps on overflow rather than returning Spark's NULL, and overflow cannot be detected at plan time, so the whole sliding decimal `SUM` case falls back. Ever-expanding frames use Comet's overflow-aware `SumDecimal` UDAF and run natively, matching Spark including overflow-to-NULL. Bigint sliding `SUM` overflow matches Spark in legacy mode (both wrap) and stays native.
- 3.5.9 and 4.1.3, audited 2026-09-30: Integral sliding `SUM` uses a native retractable accumulator in ANSI and TRY mode. It checks ordered prefix sums, so intermediate overflow throws or returns NULL even when later values cancel it, and results recover once the offending values leave the frame. Legacy sliding sums retain the DataFusion path.

[Spark Expression Support]: ../../user-guide/latest/expressions.md
22 changes: 11 additions & 11 deletions native/core/src/execution/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2480,10 +2480,9 @@ impl PhysicalPlanner {
// `PERCENT_RANK` / `CUME_DIST` / `NTILE`
// (`!uses_bounded_memory()` — "Can not execute X in a streaming
// fashion") and keeps the Spark-compatible Comet UDAFs
// (`SumDecimal` / `SumInteger` / `AvgDecimal` / `Avg`) on the
// non-streaming path since they don't implement `retract_batch`.
// Because `process_agg_func` already picks DataFusion's
// retract-capable built-ins for sliding aggregate frames,
// (`SumDecimal` / `AvgDecimal` / `Avg`) on the non-streaming path
// where needed. `process_agg_func` picks retract-capable UDAFs
// for sliding aggregate frames;
// ever-expanding aggregate frames (all that route to
// `BoundedWindowAggExec` as `PlainAggregateWindowExpr`) never
// trigger a retract call.
Expand Down Expand Up @@ -3321,9 +3320,9 @@ impl PhysicalPlanner {
// DataFusion uses `PlainAggregateWindowExpr` which does not call
// `retract_batch`, so we can safely use Comet's Spark-compatible
// UDAFs (SumDecimal/SumInteger/AvgDecimal/Avg). Otherwise it uses
// `SlidingAggregateWindowExpr` which requires retract — Comet's UDAFs
// don't implement it, so the caller must fall back to DataFusion's
// built-ins (which do).
// `SlidingAggregateWindowExpr` which requires retract. process_agg_func
// selects Comet's checked integer accumulator or a DataFusion built-in
// that supports retraction.
let is_ever_expanding = spark_expr
.spec
.as_ref()
Expand Down Expand Up @@ -3539,9 +3538,9 @@ impl PhysicalPlanner {
Some(AggExprStruct::Sum(expr)) => {
// For ever-expanding frames, use Comet's Spark-compatible Sum UDAFs
// (SumDecimal / SumInteger) which enforce Spark overflow semantics.
// For sliding frames, those UDAFs can't be used (no retract_batch),
// so delegate to DataFusion's built-in `sum`, which supports retract
// but doesn't enforce Spark's decimal precision overflow-to-NULL.
// Checked integral sliding sums also use Comet's retractable
// accumulator. Legacy sliding sums keep DataFusion's fast path;
// decimal sliding frames are rejected by the JVM planner.
let child = self.create_expr(expr.child.as_ref().unwrap(), Arc::clone(&schema))?;
let arrow_type = to_arrow_datatype(expr.datatype.as_ref().unwrap());
match arrow_type {
Expand All @@ -3556,7 +3555,8 @@ impl PhysicalPlanner {
Ok((udaf(func), vec![child]))
}
DataType::Int8 | DataType::Int16 | DataType::Int32 | DataType::Int64
if is_ever_expanding =>
if is_ever_expanding
|| from_protobuf_eval_mode(expr.eval_mode)? != EvalMode::Legacy =>
{
let eval_mode = from_protobuf_eval_mode(expr.eval_mode)?;
let func = SumInteger::try_new(arrow_type, eval_mode)?;
Expand Down
Loading
Loading