Skip to content

[SPARK-55569][SQL] Collapse non-matching pivot column values in fast-path firstAgg - #59149

Open
shrirangmhalgi wants to merge 1 commit into
apache:masterfrom
shrirangmhalgi:SPARK-55569-pivot-filter-pushdown
Open

shrirangmhalgi wants to merge 1 commit into
apache:masterfrom
shrirangmhalgi:SPARK-55569-pivot-filter-pushdown

Conversation

@shrirangmhalgi

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

In the PivotFirst fast path, firstAgg groups by (groupByKeys, pivotColumn). When the table has many distinct pivot column values outside the explicit IN list, firstAgg creates one group per distinct non-matching value - all of which
PivotFirst ignores in secondAgg.

This PR wraps the pivot column group-by key in IF(pivotColumn IN (pivotValues), pivotColumn, null), collapsing all non-matching values into a single null group. PivotFirst skips the null group as before, and the group-by keys are preserved in secondAgg, keeping query results identical to the un-optimized plan.

The optimization is skipped when:

  • The pivot column contains an aggregate expression (SPARK-24722)
  • A NULL pivot value is present (non-matching rows mapped to null would merge with legitimate null-key rows)

Why are the changes needed?

For a table with 10,000 distinct course values pivoting on 2:

  • Before: firstAgg creates 10,000 groups per group-by key (99.98% wasted)
  • After: firstAgg creates 3 groups per group-by key (2 matching + 1 null)

This reduces memory (fewer HashAggregate entries), CPU (fewer hash/compare/aggregate operations), and spill risk.

Does this PR introduce any user-facing change?

No. Query results are identical. This only affects the internal execution plan for the PivotFirst fast path.

How was this patch tested?

Five new tests in DataFramePivotSuite:

  • Plan-level assertion that If(In(...), col, null) appears in the analyzed plan
  • NULL pivot value correctly skips the optimization
  • Aggregate pivot column (SPARK-24722) correctly skips the optimization
  • Multiple aggregates (sum + avg) with the optimization
  • Group-by key with only non-pivot values is preserved as Row(year, null, null)

Existing golden file tests regenerated (pivot.sql, udf-pivot.sql, pipe-operators.sql). 49/49 DataFramePivotSuite tests pass.

Was this patch authored or co-authored using generative AI tooling?

Yes. Co-authored using Claude Opus 4.8.

…path firstAgg

In the PivotFirst fast path, the firstAgg groups by (groupByKeys, pivotColumn).
When the table has many distinct values outside the explicit pivot IN list, firstAgg
creates one group per distinct non-matching value -- all of which PivotFirst ignores
in secondAgg.

This change wraps the pivot column group-by key in
IF(pivotColumn IN (pivotValues), pivotColumn, null), collapsing all non-matching
values into a single null group. PivotFirst skips the null group as before, but
the group-by keys are preserved in secondAgg, keeping semantics identical to the
un-optimized plan.

The optimization is skipped when:
- The pivot column contains an aggregate expression (SPARK-24722)
- A NULL pivot value is present (non-matching rows would merge with legitimate
  null-key rows)

Co-authored-by: Claude Opus 4.8

@shrirangmhalgi shrirangmhalgi left a comment •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@cloud-fan / @peter-toth could you please take a look? This optimizes the PivotFirst fast path by collapsing non-matching pivot column values into a single null group in firstAgg, reducing the number of wasted groups from O(distinct_values) to O(1) per group-by key. The query results are unchanged.

This branch has not been deployed

No deployments
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