Substrait consumer ignores AggregateFunction.phase, so an intermediate aggregate runs as a complete one
- Dominant language
- Rust
- Stars
- 9.3k
- Forks
- 2.4k
- Avg merge
- 3d 7h
- Merged PRs (30d)
- 344
Description
### Describe the bug
`AggregateFunction.phase` is never read: the string `phase` does not appear anywhere under `datafusion/substrait/src/logical_plan/consumer/`. A measure declared `INITIAL_TO_INTERMEDIATE` is planned as an ordinary aggregate, and the plan is accepted with no error.
Measured on `3266eaa91`. The plans read a named table `t_avg(c0 i64 NOT NULL)` holding rows `1` and `2`, and carry one `avg` measure over `c0`:
| plan | phase | declared output | DataFusion returns |
| --- | --- | --- | --- |
| `PlanRel.rel` | `INITIAL_TO_INTERMEDIATE` | `STRUCT` | `avg(t_avg.c0):Float64?`, `1.5` |
| `PlanRel.rel` | `INITIAL_TO_RESULT` | `i64?` | `avg(t_avg.c0):Float64?`, `1.5` |
| `PlanRel.root` | `INITIAL_TO_INTERMEDIATE` | `STRUCT` | rejected: `Names list must match exactly to nested schema, but found 1 uses for 3 names` |
Two different phases, one answer — and that answer is neither declared type. `functions_arithmetic.yaml` gives `avg:i64` `return: i64?` and says it truncates for integral types, so `1.5` is DataFusion's own `avg`: not the intermediate struct the first plan asks for, and not the final value the second one declares.
The third row is why this stays out of sight. A struct-returning measure needs three names depth-first in `RelRoot` — the column, then the struct's two fields — and that plan is rejected on names, so anyone writing the rooted form reads an error about something else.
### Expected behavior
The consumer should honor the requested phase or reject it if unsupported. `LogicalPlan::Aggregate` has no phase field; partial/final aggregation is handled by the physical planner, so rejecting unsupported intermediate phases may be sufficient here.
There is also a producer compatibility issue. At the measured revision, DataFusion's producer writes phase 0 (`UNSPECIFIED`) for complete aggregates (`producer/expr/aggregate_function.rs:68`). However, [spec v0.102.0](https://github.com/substrait-io/substrait/blob/v0.102.0/proto/substrait/algebra.proto) defines `UNSPECIFIED` as `INTERMEDIATE_TO_RESULT`. Treating 0 as `INITIAL_TO_RESULT` would be a compatibility exception for those producer plans. Whether to retain that exception or update the producer should be decided explicitly.
Window functions carry the same field. I have not tested them.
### To Reproduce
The three plans are at [`probe/phase-cases`](https://github.com/alexandrefimov/substrait-conformance-cases/tree/v0.1.0/probe/phase-cases) — protobuf-JSON with the `.bin` alongside, plus the one-column table registered above. Glad to open a PR adding them as consumer tests if that is where you would want them.
Contributor guide
Research direction
Start in datafusion/substrait/src/logical_plan/consumer/ and inspect how aggregate measures are consumed, then review producer/expr/aggregate_function.rs:68. Reproduce the plans from probe/phase-cases and add consumer tests if appropriate. Done means intermediate phases are honored or explicitly rejected, with the phase-0 producer compatibility decision documented in tests or issue discussion.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- rust
- Domain
- databases
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Active
- Clarity
- Mostly clear
- Newbie friendliness
- 48/100