JoeryH opened a new issue, #24777: URL: https://github.com/apache/datafusion/issues/24777
### Describe the bug `datafusion-proto` decodes Projection, Filter, Window, Aggregate and Sort through `LogicalPlanBuilder` ([logical_plan/mod.rs](https://github.com/apache/datafusion/blob/8332cfafb37aa9209eacdbb098afc3e4bcd7f59a/datafusion/proto/src/logical_plan/mod.rs#L550-L580)), so every deserialization re-normalizes a plan that was already normalized before it was serialized. On wide plans that re-normalization is not just wasted, it is superlinear, so deserializing a cached plan gets disproportionately expensive exactly where caching a plan is most worth doing. The re-normalization work is per-column but depends only on the input plan. [`LogicalPlanBuilder::normalize`](https://github.com/apache/datafusion/blob/8332cfafb37aa9209eacdbb098afc3e4bcd7f59a/datafusion/expr/src/logical_plan/builder.rs#L1010-L1023) short-circuits for qualified columns, but for an unqualified one it recomputes, for each column: ```rust let fallback_schemas = plan.fallback_normalize_schemas(); let using_columns = plan.using_columns()?; // full apply_with_subqueries traversal ``` and [`normalize_cols`](https://github.com/apache/datafusion/blob/8332cfafb37aa9209eacdbb098afc3e4bcd7f59a/datafusion/expr/src/expr_rewriter/mod.rs#L117-L125) calls it once per column. ### To Reproduce All numbers on main @ `8332cfafb`, `--release`, M-series mac. Plan under test is `SELECT b0.. FROM (SELECT c0+1 AS a0, .. FROM t)` over an `n`-column table; the inner aliases make the outer projection's references unqualified, which is the case `normalize` cannot short-circuit. **1. End-to-end `logical_plan_from_bytes` timing.** Serialized size grows ~21x from n=100 to n=2000; decode time grows ~81x: ``` unqualified n= 100 decode 1.2ms (8576 bytes) unqualified n= 250 decode 3.2ms (22227 bytes) unqualified n= 500 decode 10.0ms (44979 bytes) unqualified n= 1000 decode 33.8ms (90479 bytes) unqualified n= 2000 decode 106.7ms (186482 bytes) ``` **2. Same plan shape, qualified outer references** (`... FROM (...) s`, projecting `s.a0 + 1`), so `normalize` short-circuits. Everything else is identical: ``` qualified n= 100 decode 946.2µs (9090 bytes) qualified n= 250 decode 2.6ms (23491 bytes) qualified n= 500 decode 7.0ms (47495 bytes) qualified n= 1000 decode 21.1ms (95495 bytes) qualified n= 2000 decode 61.2ms (196498 bytes) ``` So the normalization path alone accounts for ~43% of decode time at n=2000. **3. Isolating the builder overhead.** Taking an already-normalized wide projection and rebuilding it both ways — `LogicalPlanBuilder::from(input).project(exprs)` (what decode does) versus `Projection::try_new(exprs, input)` (what it could do) — with identical resulting schemas: ``` n= 100 builder 544.7µs ctor 307.7µs ratio 1.8x n= 500 builder 6.0ms ctor 2.2ms ratio 2.7x n= 1000 builder 21.5ms ctor 7.4ms ratio 2.9x n= 2000 builder 71.9ms ctor 17.3ms ratio 4.1x ``` The ratio widening with `n` is the point: the builder is not adding a constant factor, it is adding a superlinear term to work that is redundant for an already-normalized plan. To be clear about what these numbers do *not* say: the `ctor` column is superlinear too, so plenty of the cost is the general wide-plan planning overhead being tracked in #19795 / #7698. This issue is only about the part that deserialization does not need to pay at all. Benchmark sources: [decode timing](https://gist.github.com/JoeryH/5bc84cab356176b846ce6b27ccf18dec) and [builder-vs-constructor](https://gist.github.com/JoeryH/f17157d7ce4cd088c57ad4bd6e401b6f) — happy to attach them as a PR under `datafusion/proto/benches` instead if that is more useful. ### Expected behavior Deserializing a plan that was normalized before it was serialized should not re-normalize it, and normalization work that depends only on the input plan should not be redone once per column. ### Additional context Two independent fixes, the first strictly safe: 1. **Hoist the per-plan work out of the per-column loop.** `using_columns()` and `fallback_normalize_schemas()` depend only on `plan`, so `normalize_cols` can compute them once and reuse them across columns. No behavior change, and it helps every caller of `normalize_cols`, not just `datafusion-proto`. 2. **Decode with the node constructors in `datafusion-proto`.** `Projection::try_new`, `Filter::try_new`, `Window::try_new` and `Aggregate::try_new` skip normalization entirely. This is sound only if serialized plans are always already normalized — true for anything the SQL or DataFrame planners produced, but a real behavior change for hand-built plans, so it deserves its own discussion. Related, same root cause at other call sites: - #14563 — `DataFrame::with_column`/`with_column_renamed` were slow for this reason; fixed in #14653 by not re-normalizing already-normalized columns at that call site, which left the per-column recomputation itself in place. - #14118 — discussion of `LogicalPlan::using_columns()`. - #19795 / #7698 — the broader wide-plan planning-performance work. I'm happy to put up a PR for (1). -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
