Repository navigation
fix: Set Substrait aggregation phase to INITIAL_TO_RESULT - #25146
namanjain24-sudo wants to merge 3 commits into
Conversation
26b02c3 to
d50b923
Compare
Thank you, @Xuanwo |
ae437fa to
334f6f7
Compare
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #25146 +/- ##
==========================================
- Coverage 82.48% 82.47% -0.01%
==========================================
Files 1140 1140
Lines 437341 437342 +1
Branches 437341 437342 +1
==========================================
- Hits 360721 360710 -11
- Misses 54835 54843 +8
- Partials 21785 21789 +4 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
22cf937 to
caf289f
Compare
|
@Xuanwo just rebased onto main to clear the conflict, no code changes — but that reset the approval. Could you take another quick look when you get a chance? |
## Which issue does this PR close? - Closes apache#24967. ## Rationale for this change The Substrait consumer treats explicit intermediate aggregate phases as complete calls. This can silently return final values when a plan requests intermediate state, or report an unrelated root-schema naming error. ## What changes are included in this PR? - Validate phases in aggregate and window expressions before translating their arguments. - Accept `INITIAL_TO_RESULT` and retain `UNSPECIFIED` for compatibility with existing DataFusion-produced plans. - Reject explicit intermediate phases and unknown protobuf enum values with clear errors. - Document the compatibility limitation: Substrait defines `UNSPECIFIED` as `INTERMEDIATE_TO_RESULT`, but the consumer cannot distinguish legacy DataFusion complete calls from unspecified intermediate-state calls produced elsewhere. Those calls remain accepted. The producer change in apache#25146 is complementary and is not duplicated here. ## What is the testing strategy for this PR? - Reproduced the original behavior: an `INITIAL_TO_INTERMEDIATE` average over values 1 and 2 returned 1.5 instead of intermediate state. - Tests cover supported phases, every explicit unsupported phase, rooted and unrooted aggregates, unknown binary-protobuf enum values, and actual window output. - The full Substrait integration target passes: 213 tests passed, with six existing tests ignored. - The extended workspace suite passes 11,267 Rust tests, with eight existing tests ignored, and all 511 SQL logic-test files using an explicit four-thread limit. An initial default-concurrency run failed one ordered-aggregate spill test with a test-memory-pool exhaustion; the complete limited-concurrency rerun passed without source changes. - `cargo clippy --all-targets --all-features -- -D warnings` and the complete `./dev/rust_lint.sh` pass, including strict workspace documentation checks. ## Are there any user-facing changes? Plans with explicit unsupported aggregate or window phases now fail instead of being interpreted as complete calls. Existing unspecified-phase plans remain accepted, including the ambiguity described above. No public Rust API changes are included. Intermediate-state execution and the separate AVG output-type mismatch are not addressed here.
## Which issue does this PR close? - Closes apache#24967. ## Rationale for this change The Substrait consumer treats explicit intermediate aggregate phases as complete calls. This can silently return final values when a plan requests intermediate state, or report an unrelated root-schema naming error. ## What changes are included in this PR? - Validate phases in aggregate and window expressions before translating their arguments. - Accept `INITIAL_TO_RESULT` and retain `UNSPECIFIED` for compatibility with existing DataFusion-produced plans. - Reject explicit intermediate phases and unknown protobuf enum values with clear errors. - Document the compatibility limitation: Substrait defines `UNSPECIFIED` as `INTERMEDIATE_TO_RESULT`, but the consumer cannot distinguish legacy DataFusion complete calls from unspecified intermediate-state calls produced elsewhere. Those calls remain accepted. The producer change in apache#25146 is complementary and is not duplicated here. ## What is the testing strategy for this PR? - Reproduced the original behavior: an `INITIAL_TO_INTERMEDIATE` average over values 1 and 2 returned 1.5 instead of intermediate state. - Tests cover supported phases, every explicit unsupported phase, rooted and unrooted aggregates, unknown binary-protobuf enum values, and actual window output. - The full Substrait integration target passes: 213 tests passed, with six existing tests ignored. - The extended workspace suite passes 11,267 Rust tests, with eight existing tests ignored, and all 511 SQL logic-test files using an explicit four-thread limit. An initial default-concurrency run failed one ordered-aggregate spill test with a test-memory-pool exhaustion; the complete limited-concurrency rerun passed without source changes. - `cargo clippy --all-targets --all-features -- -D warnings` and the complete `./dev/rust_lint.sh` pass, including strict workspace documentation checks. ## Are there any user-facing changes? Plans with explicit unsupported aggregate or window phases now fail instead of being interpreted as complete calls. Existing unspecified-phase plans remain accepted, including the ambiguity described above. No public Rust API changes are included. Intermediate-state execution and the separate AVG output-type mismatch are not addressed here.
|
@Xuanwo gentle bump on this whenever you get a chance — still just the rebase, no code changes since your approval. |
The producer left `phase` at its default on every aggregate and window function call it emits, so each one carried AGGREGATION_PHASE_UNSPECIFIED. That is not the same as omitting the field: the spec gives the value the meaning INTERMEDIATE_TO_RESULT, i.e. that the arguments are already intermediate aggregation state to be combined. A LogicalPlan::Aggregate is always a complete aggregation over its input rows. The partial/final split is a physical planning concern that the logical producer has no notion of. Both `AggregateFunction.phase` and `Expression.WindowFunction.phase` are documented as required, and as needing INITIAL_TO_RESULT for a complete invocation, so set that at both call sites. This stays invisible to a DataFusion-to-DataFusion round trip because the consumer never reads `phase`, but a consumer that honours the declaration reads a complete aggregation as one whose arguments are partial state.
A DataFusion round trip cannot see a field the consumer never reads, so nothing here catches a wrong AggregationPhase: the consumer rebuilds a complete aggregation whatever the plan says. substrait-spark does read it, and maps the default UNSPECIFIED to Spark's Final, which Spark then rejects because the inputs are raw rows rather than partial buffers. An ignored Rust test writes the plan for count/sum/avg, and a small Maven project converts it with substrait-java and runs it in Spark over t(i) = 1, 2, 3, checking both the aggregation mode and the rows. Two workarounds are applied on the Java side, each named after the issue it stands in for: apache#11545 for the extension URN, and apache#25049 for the unset output_type, which can go once that fix lands. The job runs only when the Substrait crate changes, because it pulls a Spark sized dependency set.
The job runs on a plain runner rather than the amd64/rust container, so setup-builder's apt-get calls fail there. Install the protobuf compiler and add rustfmt, which is what the Substrait build needs from it.
caf289f to
6577d30
Compare
|
@Xuanwo rebased onto main again to clear the conflict (main independently added Substrait output_type for window functions, which this PR's window-phase change also touches) — no further code changes, just the rebase. Sorry for the repeat ping. |
Which issue does this PR close?
Rationale for this change
The Substrait producer never sets
phaseon the aggregate and window functioncalls it emits, so every call carries
AGGREGATION_PHASE_UNSPECIFIED.That is not the same as leaving the field out. The spec gives the value a
meaning, and it is not the one these plans need:
Both
AggregateFunction.phaseandExpression.WindowFunction.phasearedocumented as
Required. Must be set to INITIAL_TO_RESULT for ... that are not decomposable.A
LogicalPlan::Aggregateis always a complete aggregation over its inputrows. The partial/final split is a physical planning concern, and the logical
producer has no notion of it, so
INITIAL_TO_RESULTis the phase these plansshould declare. What they declare instead carries the spec meaning
INTERMEDIATE_TO_RESULT: that the arguments are already intermediate state tobe combined.
This stays invisible to a DataFusion-to-DataFusion round trip because the
consumer never reads the field (#24967). A consumer that does honour the
declaration reads a complete aggregation as one whose arguments are already
partial state.
What changes are included in this PR?
Set
phasetoAGGREGATION_PHASE_INITIAL_TO_RESULTat the two producer callsites that emit it:
from_aggregate_functioninproducer/expr/aggregate_function.rsmake_substrait_window_functioninproducer/expr/window_function.rsNo other producer site emits the field, and the consumer does not read it, so
nothing else changes.
What is the testing strategy for this PR?
New test
aggregate_and_window_functions_declare_initial_to_resultindatafusion/substrait/tests/cases/serialize.rs. It produces plans for a bareaggregate, a grouped aggregate, and a window function, then walks the produced
protobuf directly and asserts the phase on every aggregate and window function
call, rather than round-tripping through the consumer (which ignores the
field, so a round trip could not catch this).
Verified the test fails without the producer change:
The java interop test
mainand this branch both pass every test in the repository, because theDataFusion consumer never reads
phaseand rebuilds a complete aggregationeither way. #25100 asked for that gap to be covered against another
implementation, so this PR also adds one:
tests/cases/java_interop.rswrites the plan forSELECT count(i), sum(i), avg(i) FROM t. It is#[ignore]d, socargo testis unchanged for anyone without a JVM.
java-interop/converts that plan with substrait-java 0.103.0 and runs it inSpark 3.5.4 over
t(i) = 1, 2, 3, checking the aggregation mode and the rows.Run both halves with:
Against a plan from this branch it passes, with
Completeandcount 3, sum 6, avg 2.0. Against a plan frommainit fails withexpected: <[Complete, Complete, Complete]> but was: <[Final, Final, Final]>,which is the failure this change prevents.
Two workarounds are applied on the Java side, each named after the issue it
stands in for: #11545 for the extension URN, and #25049 for the unset
output_type, which can be removed once #25090 lands.The CI job runs only when
datafusion/substrait/**orrust.ymlchanges,because it pulls a Spark sized dependency set.
Are there any user-facing changes?
Plans produced by
to_substrait_plannow declareAGGREGATION_PHASE_INITIAL_TO_RESULTinstead ofAGGREGATION_PHASE_UNSPECIFIEDon aggregate and window function calls. Thisis a fix to the emitted protobuf rather than a Rust API change. Consumers that
read
phasewill now see the correct declaration; consumers that ignore it,including DataFusion's own, are unaffected.