Skip to content

fix(spark): preserve aggregate outputs in projections - #1263

Draft
bvolpato wants to merge 1 commit into
substrait-io:mainfrom
bvolpato:bvolpato/fix-spark-aggregate-project
Draft

fix(spark): preserve aggregate outputs in projections#1263
bvolpato wants to merge 1 commit into
substrait-io:mainfrom
bvolpato:bvolpato/fix-spark-aggregate-project

Conversation

@bvolpato

@bvolpato bvolpato commented Sep 3, 2026

Copy link
Copy Markdown
Member

When a Substrait Project follows an Aggregate, Spark conversion folds the projection into the aggregate but replaces the inherited outputs with only the appended expressions and ignores emit mapping. An aggregate row [1, 30] followed by a project adding 99 therefore becomes [99] instead of [1, 30, 99]; an emit selecting the original grouping column also returns the wrong value.

Project.builder()
    .input(aggregate)
    .addExpressions(ExpressionCreator.i32(false, 99))
    .remap(Rel.Remap.of(List.of(0)))
    .build();

Preserve the aggregate's existing named expressions, append the project expressions, and apply emit mapping to that complete list, as required by spec v0.102.0. Keep the existing aggregate folding so projected aggregates retain their Spark round-trip shape.

Retain inherited aggregate expressions before appending project outputs
and applying emit mapping, preserving the existing aggregate folding.
@bvolpato
bvolpato force-pushed the bvolpato/fix-spark-aggregate-project branch from 290cd5d to 35632b7 Compare September 4, 2026 16:24
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant