This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-24389-bb45fb8275bf6370923ac40bfb3dd5712b640e97 in repository https://gitbox.apache.org/repos/asf/datafusion.git
commit 47794ff2462e66ee28567ceee49a246056077bb0 Author: Reid Kaufmann <[email protected]> AuthorDate: Sat Aug 15 13:11:31 2026 +0000 Improve planning speed: Fast path for `union_schema` when all children share a schema (#24389) ## Which issue does this close? Complements #19792. Fits with the wide-`UnionExec` planning-cost work, but originates from an InfluxDB issue. ## Rationale for this change `union_schema` builds the output schema for `UnionExec` and `InterleaveExec` by coercing field metadata and nullability across **every** child. That merge is quadratic in the number of children: for each output field it walks all inputs, and for each input it walks *every other* input to union field-level metadata. For a union of `n` children with `f` fields the construction cost is `O(n^2 * f)` (worse when fields carry metadata). For narrow unions this is insignificant. It matters when a plan fans a single source out into many identical-schema children and unions them back together -- e.g. a union assembled from repartitioned copies of the same input. An instance like this occurred with InfluxDB: every child schema was the same, so the merge, guaranteed to reproduce the first child's schema, unnecessarily incurred the planning latency penalty from `O(n^2 * f)` complexity. ### Relationship to #19792 `UnionExec` construction has two quadratic halves: - **`with_new_children` / `PlanProperties`** -- addressed by #19792 (`with_new_children_and_same_properties`, `Arc<PlanProperties>`, the properties fast path). Already on `main`. - **`union_schema`** -- *not* covered by #19792 and still quadratic on `main`. This PR complements it by making `union_schema` skip the merge when it can't change the result. It deliberately doesn't touch `with_new_children`; that path is already handled. ## What changes are included? A fast path at the top of `union_schema`: after taking `inputs[0].schema()`, if every remaining child's schema is either the **same allocation** (`Arc::ptr_eq`) or **structurally equal** (`==`) to the first, return the first schema immediately. Otherwise we fall through to the existing full merge, so behavior for genuinely heterogeneous unions is byte-for-byte unchanged. ```rust let first_schema = inputs[0].schema(); if inputs[1..].iter().all(|input| { let schema = input.schema(); Arc::ptr_eq(&schema, &first_schema) || schema == first_schema }) { return Ok(first_schema); } ``` `InterleaveExec` shares `union_schema`, so it gets the same speedup for free. ## On the cost of the deep compare... The natural objection (which came up before this PR): **doesn't the deep `==` make the *unequal* case slower?** I'll paraphrase the prior conclusions, risking verbosity to avoid rehashing the discussion. Spoiler: it's not an issue. - The equal case avoids the merge, and its check is cheap. The shared-`Arc` case is settled by pointer comparison. The distinct-but-equal case runs `Schema::eq`, which is allocation-free and short-circuits on the first difference. Benchmarks show a small loss versus a pointer-equality-only control (the theoretical floor) but it still beats the full merge by a wide margin, and that advantage grows with schema complexity. - The adversarial worst case is bounded. The one shape where the scan is pure overhead is `last_differs`: children `0..n-1` are equal and the last diverges, so we scan `n` schemas, fail on the last, then merge anyway. That's a single linear `==` pass bounded by the merge that follows -- a constant fraction, not another factor of `n` -- and it takes *thousands* of near-identical children differing only in the last to hit. - Ordinary unequal unions fail fast. `SELECT a ... UNION ALL SELECT b ...` differs at field 0, so `==` rejects on the first field (see `names_differ`). And `UnionExec::try_new` already rejects misaligned children, so the only divergence `union_schema` ever sees is top-level (caught in the first pass). ## Benchmark results New bench `datafusion/physical-plan/benches/union_schema.rs` measures `UnionExec::try_new` construction over a flat schema and a nested/struct schema, for the four child shapes above. Run interleaved (baseline / patched alternated per cell) on a fixed-clock T2D VM. ### `union_schema` construction (lower is better) | scenario | n | baseline (ms) | patched (ms) | change | |---|---|---|---|---| | union_exec_try_new/shared_arc | 100 | 0.263 | 0.042 | 6.2× | | union_exec_try_new/shared_arc | 1000 | 2.60 | 0.429 | 6.1× | | union_exec_try_new/shared_arc | 4000 | 10.5 | 1.74 | 6.0× | | union_exec_try_new/content_equal | 100 | 0.262 | 0.042 | 6.2× | | union_exec_try_new/content_equal | 1000 | 2.61 | 0.433 | 6.0× | | union_exec_try_new/content_equal | 4000 | 10.5 | 1.74 | 6.0× | | union_exec_try_new/last_differs | 4000 | ~98 | ~97 | flat (±1.5%) | | union_exec_try_new/names_differ | 4000 | ~87 | ~87 | flat (±1%) | | union_exec_try_new_nested/content_equal | 1000 | 2.68 | 0.431 | 6.2× | | union_exec_try_new_nested/content_equal | 4000 | 10.8 | 1.73 | 6.2× | The `last_differs` (adversarial: N-1 children equal, deep compare then full merge) and `names_differ` (typical unequal: fails on the first field) cells were re-measured with tight per-cell interleaving (baseline/patched adjacent, 4 rounds) to control for variance: both are within ±1.5%, straddling zero. Interpretation: the **deep compare cost isn't observable end to end**. ### End-to-end planning: no regression (`sql_planner`) `cargo bench --bench sql_planner` (TPC-H + ClickBench) run baseline vs patched on the fixed-clock T2D VM. Every case lands within ±1% -- run-to-run noise -- with no case regressing beyond that noise. Notable rows, including the union-heavy `sorted_union` cases the fast path is meant to help: | case | baseline (ms) | patched (ms) | |---|---|---| | physical_plan_tpcds_all | 995.9 ± 1.4 | 991.9 ± 1.9 | | physical_plan_tpch_all | 60.4 ± 0.2 | 60.3 ± 0.1 | | physical_sorted_union_order_by_50 | 349.9 ± 2.4 | 346.3 ± 2.7 | | physical_sorted_union_order_by_10 | 12.3 ± 0.03 | 12.2 ± 0.07 | | physical_select_all_from_1000 | 30.8 ± 0.25 | 30.7 ± 0.08 | The full TPC-H q1-q22 and ClickBench sets are all flat (ratio is 1.00-1.01 in both directions). Separately, interleaving tests per-cell (baseline and patched back-to-back, so run-to-run variance -- e.g. thermal -- cancels rather than favoring one) for `physical_join_distinct` + eight ClickBench queries (4 rounds) confirmed the same thing: patched and baseline straddle zero; no systematic regression from the deep compare. ## Testing - `cargo test -p datafusion-physical-plan --lib union` -- all pass, including a new `test_union_schema_fast_path_content_equal` that exercises the `==` branch with pointer-distinct-but-equal schemas and asserts the result matches the shared schema (i.e. identical to the slow-path merge). - `cargo clippy -p datafusion-physical-plan --lib -- -D warnings` -- clean. - `cargo bench --bench union_schema` -- compiles and runs. ## Are there any user-facing changes? No: planning-time performance change only, results and schema are identical. --------- Signed-off-by: Reid Kaufmann <[email protected]> --- datafusion/physical-plan/Cargo.toml | 4 + datafusion/physical-plan/benches/union_schema.rs | 159 +++++++++++++++++++++++ datafusion/physical-plan/src/union.rs | 46 +++++++ 3 files changed, 209 insertions(+) diff --git a/datafusion/physical-plan/Cargo.toml b/datafusion/physical-plan/Cargo.toml index 0f72b74840..552b979c78 100644 --- a/datafusion/physical-plan/Cargo.toml +++ b/datafusion/physical-plan/Cargo.toml @@ -108,6 +108,10 @@ tokio = { workspace = true, features = [ harness = false name = "partial_ordering" +[[bench]] +harness = false +name = "union_schema" + [[bench]] harness = false name = "spill_io" diff --git a/datafusion/physical-plan/benches/union_schema.rs b/datafusion/physical-plan/benches/union_schema.rs new file mode 100644 index 0000000000..8cffe47466 --- /dev/null +++ b/datafusion/physical-plan/benches/union_schema.rs @@ -0,0 +1,159 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Benchmark for `UnionExec` construction cost as a function of child count. +//! +//! Scenarios (run against a flat and a nested/struct schema): +//! - `shared_arc`: every child returns the same `Arc<Schema>` +//! - `content_equal`: pointer-distinct but identical schemas per child +//! - `last_differs`: all children equal except extra metadata on the last +//! child's last top-level field. Worst case for the equality fast path: +//! the scan runs, fails on the final child, then the full merge runs anyway. +//! - `names_differ`: the first field is renamed in every child but the first, +//! the common real-world unequal union where equality fails immediately. + +use std::collections::HashMap; +use std::sync::Arc; + +use arrow::datatypes::{DataType, Field, Schema, SchemaRef}; +use criterion::{BenchmarkId, Criterion, criterion_group, criterion_main}; +use datafusion_physical_plan::ExecutionPlan; +use datafusion_physical_plan::empty::EmptyExec; +use datafusion_physical_plan::union::UnionExec; + +const NUM_FIELDS: usize = 10; +const NESTED_CHILDREN: usize = 5; +const METADATA_PER_FIELD: usize = 2; + +fn metadata(tag: &str, i: usize) -> HashMap<String, String> { + (0..METADATA_PER_FIELD) + .map(|m| (format!("key_{m}"), format!("value_{tag}_{i}_{m}"))) + .collect() +} + +fn flat_schema() -> Schema { + let fields: Vec<Field> = (0..NUM_FIELDS) + .map(|i| { + Field::new(format!("col_{i}"), DataType::Int64, true) + .with_metadata(metadata("f", i)) + }) + .collect(); + Schema::new(fields) +} + +fn nested_schema() -> Schema { + let fields: Vec<Field> = (0..NUM_FIELDS) + .map(|i| { + let children: Vec<Field> = (0..NESTED_CHILDREN) + .map(|c| { + Field::new(format!("sub_{i}_{c}"), DataType::Int64, true) + .with_metadata(metadata("n", i * NESTED_CHILDREN + c)) + }) + .collect(); + Field::new(format!("col_{i}"), DataType::Struct(children.into()), true) + .with_metadata(metadata("s", i)) + }) + .collect(); + Schema::new(fields) +} + +/// Clone `schema` with extra metadata on its last field. Divergence must stay +/// at the top level: differing nested fields change the field `DataType` +/// itself, which `UnionExec::try_new` rejects ("Schemas have to be aligned"). +fn divergent(schema: &Schema) -> Schema { + let mut fields: Vec<Field> = + schema.fields().iter().map(|f| f.as_ref().clone()).collect(); + let last = fields.pop().unwrap(); + let mut md = last.metadata().clone(); + md.insert("divergent".to_string(), "true".to_string()); + fields.push(last.with_metadata(md)); + Schema::new(fields) +} + +fn child(schema: SchemaRef) -> Arc<dyn ExecutionPlan> { + Arc::new(EmptyExec::new(schema)) +} + +fn children_shared_arc(schema: &Schema, n: usize) -> Vec<Arc<dyn ExecutionPlan>> { + let schema: SchemaRef = Arc::new(schema.clone()); + (0..n).map(|_| child(Arc::clone(&schema))).collect() +} + +fn children_content_equal(schema: &Schema, n: usize) -> Vec<Arc<dyn ExecutionPlan>> { + (0..n).map(|_| child(Arc::new(schema.clone()))).collect() +} + +/// First child keeps `schema`; the rest rename the first field — the common +/// real-world unequal union (`SELECT a .. UNION ALL SELECT b ..`), where any +/// equality check fails immediately. +fn children_names_differ(schema: &Schema, n: usize) -> Vec<Arc<dyn ExecutionPlan>> { + let mut fields: Vec<Field> = + schema.fields().iter().map(|f| f.as_ref().clone()).collect(); + let first = fields.remove(0); + let renamed = first.clone().with_name(format!("renamed_{}", first.name())); + fields.insert(0, renamed); + let alt = Schema::new(fields); + let mut children = vec![child(Arc::new(schema.clone()))]; + children.extend((1..n).map(|_| child(Arc::new(alt.clone())))); + children +} + +fn children_last_differs(schema: &Schema, n: usize) -> Vec<Arc<dyn ExecutionPlan>> { + let mut children = children_content_equal(schema, n - 1); + children.push(child(Arc::new(divergent(schema)))); + children +} + +fn bench_union_construction(c: &mut Criterion) { + for (suffix, schema, sizes) in [ + ("", flat_schema(), &[100usize, 1000, 4000][..]), + ("_nested", nested_schema(), &[1000, 4000][..]), + ] { + let mut group = c.benchmark_group(format!("union_exec_try_new{suffix}")); + for &n in sizes { + let shared = children_shared_arc(&schema, n); + let content = children_content_equal(&schema, n); + let differs = children_last_differs(&schema, n); + + group.bench_with_input( + BenchmarkId::new("shared_arc", n), + &shared, + |b, ch| b.iter(|| UnionExec::try_new(ch.clone()).unwrap()), + ); + group.bench_with_input( + BenchmarkId::new("content_equal", n), + &content, + |b, ch| b.iter(|| UnionExec::try_new(ch.clone()).unwrap()), + ); + group.bench_with_input( + BenchmarkId::new("last_differs", n), + &differs, + |b, ch| b.iter(|| UnionExec::try_new(ch.clone()).unwrap()), + ); + let names = children_names_differ(&schema, n); + group.bench_with_input( + BenchmarkId::new("names_differ", n), + &names, + |b, ch| b.iter(|| UnionExec::try_new(ch.clone()).unwrap()), + ); + } + group.finish(); + } +} + +criterion_group!(benches, bench_union_construction); +criterion_main!(benches); diff --git a/datafusion/physical-plan/src/union.rs b/datafusion/physical-plan/src/union.rs index c1cc5da31a..160772dc22 100644 --- a/datafusion/physical-plan/src/union.rs +++ b/datafusion/physical-plan/src/union.rs @@ -916,6 +916,22 @@ fn union_schema(inputs: &[Arc<dyn ExecutionPlan>]) -> Result<SchemaRef> { } let first_schema = inputs[0].schema(); + + // Fast path: when every input already shares the first input's schema, the + // field-by-field metadata/nullability merge below is redundant work that + // scales as O(n^2 * fields). This is common in practice: unions built from + // repartitioned copies of the same plan (e.g. observed in InfluxDB) hand us + // children that all carry the exact same schema. A pointer-equality check + // catches the shared-`Arc` case for free, and a content `==` comparison + // catches distinct-but-equal schemas; both let us return early and hand back + // the first schema unchanged. + if inputs[1..].iter().all(|input| { + let schema = input.schema(); + Arc::ptr_eq(&schema, &first_schema) || schema == first_schema + }) { + return Ok(first_schema); + } + let first_field_count = first_schema.fields().len(); // validate that all inputs have the same number of fields @@ -1627,6 +1643,36 @@ mod tests { Ok(()) } + #[test] + fn test_union_schema_fast_path_content_equal() -> Result<()> { + // Inputs whose schemas are pointer-distinct but structurally equal must + // take the content-equality (`==`) fast path and still produce a schema + // equal to the shared one, matching the slow-path merge exactly. + let schema = create_test_schema()?; + let distinct: SchemaRef = Arc::new((*schema).clone()); + // Guard the branch under test: these must NOT be the same allocation, so + // the fast path is reached via `==` rather than `Arc::ptr_eq`. + assert!(!Arc::ptr_eq(&schema, &distinct)); + + let memory_exec1 = + Arc::new(TestMemoryExec::try_new(&[], Arc::clone(&schema), None)?); + let memory_exec2 = + Arc::new(TestMemoryExec::try_new(&[], Arc::clone(&distinct), None)?); + let memory_exec3 = + Arc::new(TestMemoryExec::try_new(&[], Arc::clone(&distinct), None)?); + + // Capture the first child's schema before it is moved into the union. + let first_input_schema = memory_exec1.schema(); + let union_plan = + UnionExec::try_new(vec![memory_exec1, memory_exec2, memory_exec3])?; + + // The fast path returns the first child's schema Arc unchanged. Assert + // pointer equality (not just `==`): a slow-path merge would build a new, + // merely-equal Schema, so only ptr-eq proves the merge was skipped. + assert!(Arc::ptr_eq(&union_plan.schema(), &first_input_schema)); + Ok(()) + } + #[test] fn test_union_schema_mismatch() { // Test that UnionExec properly rejects inputs with different field counts --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
