Repository navigation
fix: Set Substrait output_type on aggregate functions - #25090
Conversation
3cc32f3 to
3616533
Compare
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #25090 +/- ##
=======================================
Coverage 82.33% 82.34%
=======================================
Files 1137 1137
Lines 432498 432526 +28
Branches 432498 432526 +28
=======================================
+ Hits 356116 356149 +33
+ Misses 54843 54833 -10
- Partials 21539 21544 +5 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
See comment here: #25100 (comment) |
3616533 to
8d5511a
Compare
The Substrait producer exported every aggregate call with output_type: None, even though the logical plan already knows the result type. Derive the output field from the logical expression and write it to AggregateFunction.output_type, mirroring the handling already used for scalar functions. Closes apache#25049.
8d5511a to
97bdefe
Compare
|
Maybe @gabotechs or @vbarua can help review this PR or suggest someone who is more focused on the substrait code at this time. I am sorry I don't think I am going to have time to drive this along |
|
@namanjain24-sudo I see the PR fix is just three lines of code, but the quantity of text in the PR description and your following messages seems disproportionately high for just a 3 LOC fix. If the text produced as description and messages is AI generated, please, try to review it first and summarize it as much as possible in order to be respectful with reviewers time. |
|
Fair point, sorry for the noise — I'll keep it tight going forward. Trimmed the PR description down to the essentials. |
|
I see some massive comments right after the PR description. As those relevant as well?
|
|
Good call. Deleted the workflow-approval ping (stale, checks run now), the substrait-java repro, and the interop-testing tangent — folded the one relevant fact (substrait-java accepts the plan with this fix, rejects it without) into the PR description as one sentence. Left the short apology comment above yours since it's already brief. |
gabotechs
left a comment
There was a problem hiding this comment.
Thanks @namanjain24-sudo! the fix is pretty straight forward, +1 on my side
…5367) ## Which issue does this PR close? - Closes apache#25366. ## Rationale for this change The producer left `output_type` unset on window function calls and on `LIKE` / `ILIKE`, including the `not` that wraps a negated one. Substrait documents both `Expression.WindowFunction.output_type` and `Expression.ScalarFunction.output_type` as: > Must be set to the return type of the function, exactly as derived using the declaration in the extension. A consumer that reads the field rejects such a call. substrait-java 0.103.0 refuses both plans, from `ProtoTypeConverter.from:125`: | plan produced for | error before this PR | | --- | --- | | `SELECT sum(i) OVER (ORDER BY i) FROM t` | `UnsupportedOperationException: Type is not set` at `ProtoExpressionConverter.fromWindowFunction:406` | | `SELECT i FROM t WHERE CAST(i AS VARCHAR) LIKE '1%'` | `UnsupportedOperationException: Type is not set` at `ProtoExpressionConverter:201`, via `ProtoRelConverter.newFilter` | A DataFusion round trip cannot catch this: the consumer reads `output_type` only in `consumer/expr/cast.rs`, so producer and consumer agree on the omission. apache#15831 and apache#20597 set this field for binary and unary expressions, `from_function` and higher order functions, and apache#25049 and apache#25090 cover `AggregateFunction`. These were the remaining expression kinds. ## What changes are included in this PR? Both types come from the expression itself, so they match what DataFusion derives rather than being restated in the producer: - `from_window_function` derives the type from `Expr::WindowFunction(..).to_field(schema)` and writes it on the call it builds. `make_substrait_window_function` had one caller and eight parameters once the type was added, so its body now sits in `from_window_function` and the helper is gone. - `from_like` derives the type from `Expr::Like(..).to_field(schema)` and sets it on the `like` call and on the `not` wrapper of a negated one, which yields that same type. `make_substrait_like_expr` had `from_like` as its only caller, so, as with the window helper, its body now sits in `from_like`. ## What is the testing strategy for this PR? - Two unit tests next to the existing `binary_expr_output_type`: `window_function_output_type` asserts `sum(i)` over a nullable `i64` carries a nullable `i64`, and `like_output_type` asserts a `LIKE` over a nullable input carries a nullable boolean, on the `like` call, on the `not` of a negated one, and on the `like` nested inside that `not`. - Each half was reverted on its own to check the tests pin it: without the window change only `window_function_output_type` fails, without the `LIKE` change only `like_output_type` fails. - `cargo test -p datafusion-substrait` passes (60 unit, 210 integration with the 6 that were already ignored, 3 doc tests), and `./ci/scripts/rust_clippy.sh`, the workspace clippy CI runs, is clean. - Emitted protobuf before and after, read directly rather than through a consumer: | plan | before | after | | --- | --- | --- | | window `sum` | `output_type` unset | set | | `like` | unset | set | | `not` around a negated `like` | unset | set | - With the same plans, substrait-java 0.103.0 no longer raises `Type is not set`. Both now get through `ProtoPlanConverter`, while the plans built from `main` still fail there. Each then stops further along for reasons that have nothing to do with this field: substrait-spark has no visitor for a window function that appears in a project expression, and Spark's `like` takes two arguments while the call we emit passes three, the third being the escape character. That second one is apache#25442, which emits `LIKE` in the two-argument form the extension defines. ## Are there any user-facing changes? No API changes. Plans produced for window functions and `LIKE` now carry the expression's type, which consumers that require it will accept.
## Which issue does this PR close? Closes apache#25049. ## Rationale for this change The Substrait producer exported aggregate calls with `output_type: None`. Per spec, `output_type` should carry the function's return type — this mirrors the fix already done for scalar functions in apache#20597. Verified against substrait-java 0.103.0: it rejects an aggregate measure with no `output_type` (`UnsupportedOperationException: Type is not set`) and accepts the same plan once this field is set, recovering the correct type. ## What changes are included in this PR? `from_aggregate_function` now derives the output field via `Expr::AggregateFunction(..).to_field(schema)` and sets it with `to_substrait_type_from_field`, the same path scalar functions already use. ## What is the testing strategy for this PR? Added `aggregate_function_output_type` unit test covering `count`, `sum`, `avg`, `min`. Fails on `main`, passes with this change. Existing `datafusion-substrait` suite passes unchanged. ## Are there any user-facing changes? Substrait plans now declare `output_type` on aggregate calls. No public API changes.
…5367) ## Which issue does this PR close? - Closes apache#25366. ## Rationale for this change The producer left `output_type` unset on window function calls and on `LIKE` / `ILIKE`, including the `not` that wraps a negated one. Substrait documents both `Expression.WindowFunction.output_type` and `Expression.ScalarFunction.output_type` as: > Must be set to the return type of the function, exactly as derived using the declaration in the extension. A consumer that reads the field rejects such a call. substrait-java 0.103.0 refuses both plans, from `ProtoTypeConverter.from:125`: | plan produced for | error before this PR | | --- | --- | | `SELECT sum(i) OVER (ORDER BY i) FROM t` | `UnsupportedOperationException: Type is not set` at `ProtoExpressionConverter.fromWindowFunction:406` | | `SELECT i FROM t WHERE CAST(i AS VARCHAR) LIKE '1%'` | `UnsupportedOperationException: Type is not set` at `ProtoExpressionConverter:201`, via `ProtoRelConverter.newFilter` | A DataFusion round trip cannot catch this: the consumer reads `output_type` only in `consumer/expr/cast.rs`, so producer and consumer agree on the omission. apache#15831 and apache#20597 set this field for binary and unary expressions, `from_function` and higher order functions, and apache#25049 and apache#25090 cover `AggregateFunction`. These were the remaining expression kinds. ## What changes are included in this PR? Both types come from the expression itself, so they match what DataFusion derives rather than being restated in the producer: - `from_window_function` derives the type from `Expr::WindowFunction(..).to_field(schema)` and writes it on the call it builds. `make_substrait_window_function` had one caller and eight parameters once the type was added, so its body now sits in `from_window_function` and the helper is gone. - `from_like` derives the type from `Expr::Like(..).to_field(schema)` and sets it on the `like` call and on the `not` wrapper of a negated one, which yields that same type. `make_substrait_like_expr` had `from_like` as its only caller, so, as with the window helper, its body now sits in `from_like`. ## What is the testing strategy for this PR? - Two unit tests next to the existing `binary_expr_output_type`: `window_function_output_type` asserts `sum(i)` over a nullable `i64` carries a nullable `i64`, and `like_output_type` asserts a `LIKE` over a nullable input carries a nullable boolean, on the `like` call, on the `not` of a negated one, and on the `like` nested inside that `not`. - Each half was reverted on its own to check the tests pin it: without the window change only `window_function_output_type` fails, without the `LIKE` change only `like_output_type` fails. - `cargo test -p datafusion-substrait` passes (60 unit, 210 integration with the 6 that were already ignored, 3 doc tests), and `./ci/scripts/rust_clippy.sh`, the workspace clippy CI runs, is clean. - Emitted protobuf before and after, read directly rather than through a consumer: | plan | before | after | | --- | --- | --- | | window `sum` | `output_type` unset | set | | `like` | unset | set | | `not` around a negated `like` | unset | set | - With the same plans, substrait-java 0.103.0 no longer raises `Type is not set`. Both now get through `ProtoPlanConverter`, while the plans built from `main` still fail there. Each then stops further along for reasons that have nothing to do with this field: substrait-spark has no visitor for a window function that appears in a project expression, and Spark's `like` takes two arguments while the call we emit passes three, the third being the escape character. That second one is apache#25442, which emits `LIKE` in the two-argument form the extension defines. ## Are there any user-facing changes? No API changes. Plans produced for window functions and `LIKE` now carry the expression's type, which consumers that require it will accept.
Which issue does this PR close?
Closes #25049.
Rationale for this change
The Substrait producer exported aggregate calls with
output_type: None. Per spec,output_typeshould carry the function's return type — this mirrors the fix alreadydone for scalar functions in #20597.
Verified against substrait-java 0.103.0: it rejects an aggregate measure with no
output_type(UnsupportedOperationException: Type is not set) and accepts the sameplan once this field is set, recovering the correct type.
What changes are included in this PR?
from_aggregate_functionnow derives the output field viaExpr::AggregateFunction(..).to_field(schema)and sets it withto_substrait_type_from_field, the same path scalar functions already use.What is the testing strategy for this PR?
Added
aggregate_function_output_typeunit test coveringcount,sum,avg,min.Fails on
main, passes with this change. Existingdatafusion-substraitsuite passesunchanged.
Are there any user-facing changes?
Substrait plans now declare
output_typeon aggregate calls. No public API changes.