Skip to content

perf: reuse column normalization context across expressions - #25010

Merged
alamb merged 1 commit into
apache:mainfrom
goutamadwant:fix-normalization-context-24777
Sep 20, 2026
Merged

alamb merged 1 commit into
apache:mainfrom
goutamadwant:fix-normalization-context-24777

Conversation

@goutamadwant

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Rationale for this change

Column normalization repeatedly collects fallback schemas and traverses the input plan for USING columns. Wide expression lists and projections repeat that work for the same immutable plan.

What changes are included in this PR?

  • Introduce a private, lazy normalization context reused across columns and expression lists.
  • Share the context in sort normalization and validated projection construction, including wildcard expansion.
  • Keep already-qualified columns and expressions that do not need normalization on the existing fast path.
  • Add benchmarks for qualified and unqualified expressions and projection construction at several schema widths.

What is the testing strategy for this PR?

  • Add normalize_batch_schema_precedence, normalize_batch_using_join, and normalize_batch_skips_unused_plan_context to cover schema precedence, USING joins, ambiguity/error order, sort options, and lazy handling of qualified columns and literals.
  • In balanced local release-nonlto runs, constructing a 2,000-column unqualified projection falls from about 100 ms to 35 ms. This measures projection construction, not full protobuf decoding; small controls remain noisy.
  • Reproduce with cargo bench -p datafusion-expr --bench normalize_columns --profile release-nonlto.
  • Focused expression and SQL tests pass. The required extended workspace test command also passes, including all 511 SQL logic-test files.
  • Expression-crate Clippy passes with all targets and features enabled. The complete documented dev/rust_lint.sh also passes, including strict workspace documentation checks.
  • Full-workspace Clippy with all features enabled hits the existing PostgreSQL decimal-formatting lint in Full workspace clippy fails in PostgreSQL decimal formatting #24974; the affected source is unchanged here.

Are there any user-facing changes?

No public API or name-resolution behavior changes are intended. Normalization reuses plan context instead of collecting it for each column; existing schema lookup costs remain.

@github-actions github-actions Bot added the logical-expr Logical plan and expressions label Sep 7, 2026
@codecov-commenter

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 75.45455% with 27 lines in your changes missing coverage. Please review.
✅ Project coverage is 81.70%. Comparing base (e1ca94f) to head (2fcda6c).
⚠️ Report is 2 commits behind head on main.

Files with missing lines Patch % Lines
datafusion/expr/src/expr_rewriter/mod.rs 77.14% 1 Missing and 23 partials ⚠️
datafusion/expr/src/logical_plan/builder.rs 40.00% 0 Missing and 3 partials ⚠️
Additional details and impacted files
@@           Coverage Diff            @@
##             main   #25010    +/-   ##
========================================
  Coverage   81.69%   81.70%            
========================================
  Files        1127     1127            
  Lines      415471   415743   +272     
  Branches   415471   415743   +272     
========================================
+ Hits       339424   339672   +248     
+ Misses      56108    56107     -1     
- Partials    19939    19964    +25     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@kosiew

kosiew commented Sep 16, 2026

Copy link
Copy Markdown
Contributor

benchmark run

group                                             HEAD                                   bench-25010
-----                                             ----                                   -----------
normalize_columns/expressions/qualified/10        1.02      4.2±0.03µs        ? ?/sec    1.00      4.1±0.04µs        ? ?/sec
normalize_columns/expressions/qualified/100       1.00     42.4±0.32µs        ? ?/sec    1.00     42.1±0.52µs        ? ?/sec
normalize_columns/expressions/qualified/2000      1.01    864.4±7.16µs        ? ?/sec    1.00    854.6±8.82µs        ? ?/sec
normalize_columns/expressions/qualified/500       1.01    212.1±1.85µs        ? ?/sec    1.00    209.9±2.25µs        ? ?/sec
normalize_columns/expressions/unqualified/10      1.54      6.5±0.05µs        ? ?/sec    1.00      4.2±0.04µs        ? ?/sec
normalize_columns/expressions/unqualified/100     3.37    231.2±0.51µs        ? ?/sec    1.00     68.6±0.39µs        ? ?/sec
normalize_columns/expressions/unqualified/2000    11.90    74.7±0.18ms        ? ?/sec    1.00      6.3±0.02ms        ? ?/sec
normalize_columns/expressions/unqualified/500     7.45      4.7±0.01ms        ? ?/sec    1.00    631.2±2.81µs        ? ?/sec
normalize_columns/projection/qualified/10         1.01     36.2±0.61µs        ? ?/sec    1.00     35.7±0.29µs        ? ?/sec
normalize_columns/projection/qualified/100        1.00    403.8±1.39µs        ? ?/sec    1.02    410.5±3.15µs        ? ?/sec
normalize_columns/projection/qualified/2000       1.00     26.9±0.17ms        ? ?/sec    1.06     28.5±0.29ms        ? ?/sec
normalize_columns/projection/qualified/500        1.00      3.2±0.01ms        ? ?/sec    1.03      3.3±0.01ms        ? ?/sec
normalize_columns/projection/unqualified/10       1.06     40.1±0.36µs        ? ?/sec    1.00     37.6±0.20µs        ? ?/sec
normalize_columns/projection/unqualified/100      1.39    608.7±1.86µs        ? ?/sec    1.00    437.7±1.92µs        ? ?/sec
normalize_columns/projection/unqualified/2000     3.74     94.2±0.46ms        ? ?/sec    1.00     25.2±0.07ms        ? ?/sec
normalize_columns/projection/unqualified/500      2.31      7.3±0.01ms        ? ?/sec    1.00      3.2±0.01ms        ? ?/sec

@kosiew kosiew left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@goutamadwant,

Thanks for working on this. The lazy reusable normalization context looks like a nice improvement, and I like the added correctness tests and benchmark coverage.

I left one non-blocking benchmark suggestion below.

.map(|i| Field::new(format!("c{i}"), DataType::Int32, false))
.collect::<Vec<_>>(),
);
let input = table_scan(Some("t"), &schema, None)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice to see benchmark coverage added here. The unqualified case already exercises the cached using_columns() traversal, but this plan does not include any JOIN ... USING nodes. It could be useful to add a wide JOIN ... USING input with unqualified projected expressions as well. That would measure the populated USING-column path directly and help protect that performance improvement going forward. Non-blocking.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@goutamadwant perhaps you can add this new benchmark as a follow on PR

@alamb
alamb added this pull request to the merge queue Sep 20, 2026
@alamb

alamb commented Sep 20, 2026

Copy link
Copy Markdown
Contributor

Merging this one in to try and keep the code flowing (our backlog of now approved PRs is getting quite large)

@alamb

alamb commented Sep 20, 2026

Copy link
Copy Markdown
Contributor

Thank you @goutamadwant and @kosiew

Merged via the queue into apache:main with commit 7f3e657 Sep 20, 2026
42 checks passed
rluvaton pushed a commit to rluvaton/datafusion that referenced this pull request Oct 1, 2026
…he#25739)

## Which issue does this PR close?

- Part of apache#25248.

## Rationale for this change

Planning a `SELECT` with many aggregate expressions takes time that
grows with the square of the number of expressions. The report measured
2.7 s of SQL planning for 2,500 items, and a profile shows almost all of
it goes into building the projection above the aggregate.

For each select item, the planner replaces sub-expressions that the
aggregate already computes with references to its output columns. It
rebuilt the lookup table of those expressions for every item, so N items
built N tables of N entries. Hashing an aggregate expression also hashes
its function signature, which makes each rebuild slow.

## What changes are included in this PR?

- The projection builder now builds that table once and reuses it for
every item, through a crate-private `Columnizer`. `columnize_expr` keeps
its signature and uses the same helper. This follows apache#25010, which
reuses the column normalization context in the same function.
- The `sql_planner` benchmark gains a case with 1,000 aggregates.

Very wide selects still do some quadratic work. Two SQL planner helpers,
`rebase_expr` and `find_aggregate_exprs`, compare each expression
against a list, and schema field lookups scan every field. I'd rather
handle those in separate PRs, so this one only references the issue.

## What is the testing strategy for this PR?

The rewrite gives the same result as before, so the existing expression,
SQL planner and optimizer tests and the sqllogictest suite cover it, and
they pass locally.

`sql_planner` benchmark against `main` on an 18-core M-series Mac,
`release-nonlto`:

| Benchmark | main | this PR | change |
| --- | --- | --- | --- |
| `logical_wide_aggregate_1000_exprs` (new) | 182.3 ms | 31.1 ms | -83%
|
| `logical_wide_aggregate_100_exprs` | 2.43 ms | 0.89 ms | -63% |
| `physical_select_aggregates_from_200` | 6.98 ms | 4.69 ms | -33% |

The query shape from the issue, `CAST(sum(id + i) AS VARCHAR)` repeated
N times, planned with `EXPLAIN` in `datafusion-cli` (unoptimized `ci`
build, so absolute times are high):

| N | main | this PR |
| --- | --- | --- |
| 500 | 0.72 s | 0.66 s |
| 1000 | 1.56 s | 0.87 s |
| 2000 | 5.05 s | 1.83 s |
| 4000 | 18.83 s | 4.57 s |

<details>
<summary>Commands and raw output</summary>

Benchmarks. The new case is copied onto `main` so both runs measure the
same queries. The `sql_planner` bench checks that
`benchmarks/data/hits_partitioned` exists at startup; I used a one-row
stand-in, since none of these three cases read it.

```
PR=eb8f5ce9d4d0b50cffb97648188c84a4034ec92c
git fetch origin main "$PR"
git switch --detach origin/main
git checkout "$PR" -- datafusion/core/benches/sql_planner.rs
cargo bench -p datafusion --bench sql_planner --profile release-nonlto -- \
  'logical_wide_aggregate|physical_select_aggregates_from_200' --save-baseline main
git checkout -f "$PR"
cargo bench -p datafusion --bench sql_planner --profile release-nonlto -- \
  'logical_wide_aggregate|physical_select_aggregates_from_200' --baseline main
```

`main`:

```
physical_select_aggregates_from_200
                        time:   [6.8908 ms 6.9799 ms 7.0789 ms]
logical_wide_aggregate_100_exprs
                        time:   [2.3629 ms 2.4285 ms 2.4977 ms]
logical_wide_aggregate_1000_exprs
                        time:   [180.60 ms 182.29 ms 184.28 ms]
```

This PR:

```
physical_select_aggregates_from_200
                        time:   [4.6522 ms 4.6890 ms 4.7269 ms]
                        change: [−33.854% −32.821% −31.767%] (p = 0.00 < 0.05)
                        Performance has improved.
logical_wide_aggregate_100_exprs
                        time:   [881.22 µs 888.79 µs 897.96 µs]
                        change: [−63.824% −62.709% −61.533%] (p = 0.00 < 0.05)
                        Performance has improved.
logical_wide_aggregate_1000_exprs
                        time:   [30.705 ms 31.131 ms 31.628 ms]
                        change: [−83.222% −82.922% −82.577%] (p = 0.00 < 0.05)
                        Performance has improved.
```

`datafusion-cli` timing, with the CLI built by `cargo build --profile ci
-p datafusion-cli` on each side:

```bash
#!/usr/bin/env bash
# Usage: wide-timing.sh <datafusion-cli>
set -euo pipefail
python3 - <<'EOF'
for n in (500, 1000, 2000, 4000):
    items = ', '.join(f'CAST(sum(id + {i}) AS VARCHAR) AS c{i}' for i in range(n))
    with open(f'wide_{n}.sql', 'w') as f:
        f.write('CREATE TABLE t AS SELECT value AS id FROM range(1000);\n')
        f.write(f'EXPLAIN SELECT {items} FROM t;\n')
EOF
for n in 500 1000 2000 4000; do
  /usr/bin/time -p "$1" -q -f "wide_$n.sql" 2>&1 >/dev/null | awk -v n="$n" '/^real/ {print "N=" n, $2 "s"}'
done
```

```
== main
N=500 0.72s
N=1000 1.56s
N=2000 5.05s
N=4000 18.83s
== this PR
N=500 0.66s
N=1000 0.87s
N=2000 1.83s
N=4000 4.57s
```

</details>

## Are there any user-facing changes?

No. Planning is faster, and `columnize_expr` keeps its signature and
behavior.

Signed-off-by: Stefan Wang <1fannnw@gmail.com>
Co-authored-by: Kumar Ujjawal <ujjawalpathak6@gmail.com>
rluvaton pushed a commit to rluvaton/datafusion that referenced this pull request Oct 1, 2026
## Which issue does this PR close?

Closes apache#24777

## Rationale for this change

Deserializing wide logical plans with `datafusion-proto` can take
disproportionately longer as the number of projected expressions
increases. This is because deserialization re-runs expression
normalization on logical plan nodes that have already been normalized
before serialization.

Avoiding this redundant work makes logical plan deserialization more
efficient, particularly for wide plans.

## What changes are included in this PR?

Decode `Projection`, `Filter`, `Window`, `Aggregate`, and `Sort` logical
plan nodes directly through their constructors instead of rebuilding
them through `LogicalPlanBuilder`.

This avoids the redundant expression normalization performed by the
builder methods while preserving the serialized logical plan structure.

This implements the constructor-based approach described as fix (2) in
apache#24777 and is complementary to apache#25010, which optimizes the normalization
work itself.

## What is the testing strategy for this PR?

The change is covered by the existing logical plan protobuf roundtrip
tests, which verify that the decoded plans remain equivalent to the
serialized plans.

The following tests were run successfully:

* `cargo check -p datafusion-proto`
* `cargo test -p datafusion-proto --test proto_integration
roundtrip_logical_plan`
* `cargo test -p datafusion-proto --test proto_integration
roundtrip_logical_plan_aggregation`
* `cargo test -p datafusion-proto --test proto_integration
roundtrip_logical_plan_sort`
* `cargo test -p datafusion-proto --test proto_integration
roundtrip_window`

No new tests were added because the existing protobuf logical-plan
roundtrip coverage exercises the affected decode paths.

The full `proto_integration` test suite was also run. It had 254 passing
tests and 7 failures caused by missing Parquet test data from the
`parquet-testing` submodule; these failures are unrelated to this
change.

## Are there any user-facing changes?

No.

Co-authored-by: kosiew <kosiew@gmail.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

logical-expr Logical plan and expressions v56.0.0

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants