From 7b13c4c7c686e41b07901cff2c3bdcee05e9fa88 Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Fri, 31 Jul 2026 01:02:34 -0500 Subject: [PATCH] fix: keep a CoalescePartitionsExec required by a SinglePartition child (#23948) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - None filed; happy to open one if preferred. A valid query can be planned into a physical plan that `SanityCheckPlan` then rejects: ``` SanityCheckPlan caused by Error during planning: Plan: ["HashJoinExec: mode=CollectLeft, join_type=Left, on=[(id@0, id@0)], projection=[id@0]", " DataSourceExec: file_groups={4 groups: [...]}, projection=[id], file_type=parquet", " RepartitionExec: partitioning=RoundRobinBatch(8), input_partitions=1", " CoalescePartitionsExec", " ProjectionExec: expr=[first_value(t.id) ORDER BY [...]@1 as id]", " AggregateExec: mode=FinalPartitioned, gby=[id@0 as id], aggr=[first_value(t.id) ORDER BY [...]]", " RepartitionExec: partitioning=Hash([id@0], 8), input_partitions=4", " AggregateExec: mode=Partial, gby=[id@1 as id], aggr=[first_value(t.id) ORDER BY [...]]", " DataSourceExec: file_groups={4 groups: [...]}, projection=[ts, id], file_type=parquet"] does not satisfy distribution requirements: SinglePartition. Child-0 output partitioning: UnknownPartitioning(4) ``` The `HashJoinExec` is in `CollectLeft` mode, which requires `Distribution::SinglePartition` on its build (left) child, but child 0 is a bare 4-partition `DataSourceExec` with no `CoalescePartitionsExec` above it. Self-contained reproducer with `datafusion-cli` (the four `COPY` statements are what make the scan multi-partition): ```sql set datafusion.execution.target_partitions = 8; set datafusion.optimizer.repartition_file_scans = false; create table src (id int, ts int) as values (1, 10), (2, 20), (3, 30); copy (select * from src) to 'data/0.parquet' stored as parquet; copy (select * from src) to 'data/1.parquet' stored as parquet; copy (select * from src) to 'data/2.parquet' stored as parquet; copy (select * from src) to 'data/3.parquet' stored as parquet; create external table t stored as parquet location 'data/'; select a.id from t a left join (select distinct on (id) id, ts from t order by id, ts) f on a.id = f.id order by a.id; ``` Setting `datafusion.optimizer.repartition_sorts = false` makes it plan fine, which points at the sort-parallelization phase. `EnsureRequirements` does insert the coalesce for the `SinglePartition` requirement (`enforce_distribution.rs`, `Distribution::SinglePartition => add_merge_on_top(...)`). Its own phase 3a (`parallelize_sorts`) then takes it back out: `remove_bottleneck_in_subplan` removes a `CoalescePartitionsExec` found at `children[0]` positionally, without consulting the parent's distribution requirement for that child. That parent is reached because `update_coalesce_ctx_children` marks a node as connected when *any* child qualifies. It correctly excludes a `SinglePartition`-requiring child from *setting* the flag, but the join's other child (`UnspecifiedDistribution`, connected to a coalesce below) sets it, so the traversal descends into the join and rewrites child 0 anyway. Nothing re-enforces distribution afterwards, so `SanityCheckPlan` is the first thing to notice. Note the surviving `CoalescePartitionsExec` on the probe side in the plan above: it is what propagated the flag, and it is untouched because the `if` returns without recursing into child 1. The sibling helper on the phase 2b path already does consult the requirement (`update_child_to_remove_unnecessary_sort` / `remove_corresponding_sort_from_sub_plan` re-add a merge using the per-child `child_distribution(child_idx)`); only this path is missing it. The same failure shows up with a build child that is already hash-partitioned on the join key (`Child-0 output partitioning: Hash([k@0], 8)`), which is what a `JoinSelection` input swap leaves behind — a `CollectLeft` join reported as `join_type=Right` with an embedded projection. `remove_bottleneck_in_subplan` now checks the parent's per-child distribution requirement before removing a coalesce, both for `children[0]` and when recursing into the other children. The node `parallelize_sorts` is itself rewriting (the root of the call) is exempt, since the caller drops that node and rebuilds the sort cascade around the result — that is the rule's intended transformation, and gating it too would disable sort parallelization below a global sort. This is threaded through as an `is_root` flag on a private `_impl` function; the public entry point keeps its signature. Yes, at two levels: - An end-to-end sqllogictest in `datafusion/sqllogictest/test_files/joins.slt` reproducing it from SQL (the reproducer above, with the data written by `COPY` inside the test). On `main` it fails with exactly the distribution error above. - Two tests in `datafusion/core/tests/physical_optimizer/ensure_requirements.rs` covering both shapes of the build child (`UnknownPartitioning(n)` and `Hash([k], n)`), running the full `EnsureRequirements` rule and then `SanityCheckPlan` via the existing `optimize_and_sanity_check` helper, plus the idempotency check. `cargo test -p datafusion-physical-optimizer`, `cargo test -p datafusion --test core_integration -- physical_optimizer` (530 tests) and the full `sqllogictest` suite (498 files) pass. No API changes. Plans that were previously rejected by `SanityCheckPlan` now plan and execute; a coalesce that is genuinely required is retained where it was previously (incorrectly) removed. --------- Co-authored-by: Claude Opus 5 (cherry picked from commit 6636d6be253255c03e1fd5fda8694004f45aaec3) --- .../physical_optimizer/enforce_sorting.rs | 117 +++++++++++++++++ .../src/enforce_sorting/mod.rs | 47 ++++++- datafusion/sqllogictest/test_files/joins.slt | 118 ++++++++++++++++++ 3 files changed, 278 insertions(+), 4 deletions(-) diff --git a/datafusion/core/tests/physical_optimizer/enforce_sorting.rs b/datafusion/core/tests/physical_optimizer/enforce_sorting.rs index 40bcdbbd6efef..3336caa94b604 100644 --- a/datafusion/core/tests/physical_optimizer/enforce_sorting.rs +++ b/datafusion/core/tests/physical_optimizer/enforce_sorting.rs @@ -2850,3 +2850,120 @@ async fn test_sort_with_streaming_table() -> Result<()> { Ok(()) } + +/// Builds the plan shape that `parallelize_sorts` sees for a `CollectLeft` `HashJoinExec` +/// whose build side is required to be `SinglePartition`, i.e. the output of the +/// distribution + sorting phases, not a freshly planned tree: +/// +/// ```text +/// CoalescePartitionsExec <- the node `parallelize_sorts` rewrites +/// HashJoinExec: mode=CollectLeft +/// CoalescePartitionsExec <- satisfies `SinglePartition` on the build side +/// +/// RepartitionExec: RoundRobinBatch +/// CoalescePartitionsExec <- links the join into the coalesce cascade +/// +/// ``` +/// +/// Both coalesces below the join matter. The probe-side one is what makes +/// `update_coalesce_ctx_children` mark the join as connected — it only skips children that +/// require `SinglePartition`, and the probe side does not — so the walk descends into the +/// join. The build-side one is the one that must survive. +fn collect_left_plan_before_parallelize_sorts( + build: Arc, + join_type: JoinType, +) -> Result> { + use datafusion_common::NullEquality; + use datafusion_physical_plan::joins::{HashJoinExec, PartitionMode}; + + let schema = create_test_schema2()?; + let build = coalesce_partitions_exec(build); + let probe = repartition_exec(coalesce_partitions_exec(parquet_exec(schema))); + + let on = vec![( + Arc::new(Column::new("col_a", 0)) as _, + Arc::new(Column::new("col_a", 0)) as _, + )]; + let join: Arc = Arc::new(HashJoinExec::try_new( + build, + probe, + on, + None, + &join_type, + None, + PartitionMode::CollectLeft, + NullEquality::NullEqualsNothing, + false, + )?); + + Ok(coalesce_partitions_exec(join)) +} + +/// A `CollectLeft` `HashJoinExec` requires `Distribution::SinglePartition` on its build +/// (left) child, so the distribution phase puts a `CoalescePartitionsExec` on top of a +/// multi-partition build side. The sort-parallelization phase (`parallelize_sorts`) must +/// not take that coalesce back out again. +/// +/// It used to, because `remove_bottleneck_in_subplan` removed a coalesce found at +/// `children[0]` positionally, without consulting the parent's distribution requirement for +/// that child. The result was a build side left multi-partition with nothing to re-enforce +/// distribution afterwards, which `SanityCheckPlan` rejected with "does not satisfy +/// distribution requirements: SinglePartition". +/// +/// Cherry-picked from apache/datafusion#23948. +#[test] +fn test_collect_left_join_keeps_build_side_coalesce() -> Result<()> { + let schema = create_test_schema2()?; + let build = repartition_exec(parquet_exec(schema)); + let plan = collect_left_plan_before_parallelize_sorts(build, JoinType::Left)?; + + let ctx = PlanWithCorrespondingCoalescePartitions::new_default(plan); + let rewritten = ctx.transform_up(parallelize_sorts).data()?.plan; + + // The build-side coalesce is retained; the probe-side one is still removed, which is + // the parallelization this phase exists for. + assert_snapshot!(displayable(rewritten.as_ref()).indent(true).to_string(), @r" + CoalescePartitionsExec + HashJoinExec: mode=CollectLeft, join_type=Left, on=[(col_a@0, col_a@0)] + CoalescePartitionsExec + RepartitionExec: partitioning=RoundRobinBatch(10), input_partitions=1 + DataSourceExec: file_groups={1 group: [[x]]}, projection=[col_a, col_b], file_type=parquet + RepartitionExec: partitioning=RoundRobinBatch(10), input_partitions=1 + DataSourceExec: file_groups={1 group: [[x]]}, projection=[col_a, col_b], file_type=parquet + "); + + Ok(()) +} + +/// The same removal, with a build side that is already hash-partitioned on the join key +/// rather than round-robin partitioned. This is the shape a `JoinSelection` input swap +/// leaves behind (a `CollectLeft` join reported as `join_type=Right`) when the build +/// subtree is the output of an aggregate or a partitioned join: the build side satisfies +/// the join's *hash* requirement but still not `SinglePartition`, so the coalesce is just +/// as load-bearing. +/// +/// Cherry-picked from apache/datafusion#23948. +#[test] +fn test_collect_left_join_keeps_hash_partitioned_build_side_coalesce() -> Result<()> { + let schema = create_test_schema2()?; + let build: Arc = Arc::new(RepartitionExec::try_new( + parquet_exec(schema), + Partitioning::Hash(vec![Arc::new(Column::new("col_a", 0))], 10), + )?); + let plan = collect_left_plan_before_parallelize_sorts(build, JoinType::Right)?; + + let ctx = PlanWithCorrespondingCoalescePartitions::new_default(plan); + let rewritten = ctx.transform_up(parallelize_sorts).data()?.plan; + + assert_snapshot!(displayable(rewritten.as_ref()).indent(true).to_string(), @r" + CoalescePartitionsExec + HashJoinExec: mode=CollectLeft, join_type=Right, on=[(col_a@0, col_a@0)] + CoalescePartitionsExec + RepartitionExec: partitioning=Hash([col_a@0], 10), input_partitions=1 + DataSourceExec: file_groups={1 group: [[x]]}, projection=[col_a, col_b], file_type=parquet + RepartitionExec: partitioning=RoundRobinBatch(10), input_partitions=1 + DataSourceExec: file_groups={1 group: [[x]]}, projection=[col_a, col_b], file_type=parquet + "); + + Ok(()) +} diff --git a/datafusion/physical-optimizer/src/enforce_sorting/mod.rs b/datafusion/physical-optimizer/src/enforce_sorting/mod.rs index 241c843556e1c..2adc702e6b1c1 100644 --- a/datafusion/physical-optimizer/src/enforce_sorting/mod.rs +++ b/datafusion/physical-optimizer/src/enforce_sorting/mod.rs @@ -676,11 +676,45 @@ fn adjust_window_sort_removal( /// the plan, some of the remaining `RepartitionExec`s might become unnecessary. /// Removes such `RepartitionExec`s from the plan as well. fn remove_bottleneck_in_subplan( + requirements: PlanWithCorrespondingCoalescePartitions, +) -> Result { + // The root is the node `parallelize_sorts` is rewriting (a `SortExec`, + // `SortPreservingMergeExec` or `CoalescePartitionsExec`). Its own distribution + // requirement does not constrain the removal, because the caller drops the node and + // rebuilds the cascade around the result. + remove_bottleneck_in_subplan_impl(requirements, true) +} + +fn remove_bottleneck_in_subplan_impl( mut requirements: PlanWithCorrespondingCoalescePartitions, + is_root: bool, ) -> Result { let plan = &requirements.plan; + // Below the root, a `CoalescePartitionsExec` feeding a child that requires + // `Distribution::SinglePartition` is not an avoidable bottleneck: it is what satisfies + // that requirement. Removing it leaves the parent with a multi-partition input it cannot + // accept, and nothing re-runs distribution enforcement afterwards, so the plan reaches + // `SanityCheckPlan` invalid. The traversal reaches such a node because + // `update_coalesce_ctx_children` marks a node as connected when *any* child qualifies: + // a `CollectLeft` `HashJoinExec` whose probe side is connected is descended into even + // though its build side must stay single-partition. + // + // Only `SinglePartition` is protected. A `HashPartitioned` child is in principle in the + // same position — a single-partition input trivially satisfies a hash requirement, so a + // coalesce below one is also load-bearing — but nothing puts a coalesce there: + // `ensure_distribution` satisfies a hash requirement with a `RepartitionExec`, never a + // `CoalescePartitionsExec`. Widening the check would be dead code today. + let dist_reqs = plan.required_input_distribution(); + let removable = |idx: usize| { + is_root || !matches!(dist_reqs.get(idx), Some(Distribution::SinglePartition)) + }; + let remove_from_first_child = requirements + .children + .first() + .is_some_and(|child| is_coalesce_partitions(&child.plan)) + && removable(0); let children = &mut requirements.children; - if is_coalesce_partitions(&children[0].plan) { + if remove_from_first_child { // We can safely use the 0th index since we have a `CoalescePartitionsExec`. let mut new_child_node = children[0].children.swap_remove(0); while new_child_node.plan.output_partitioning() == plan.output_partitioning() @@ -694,9 +728,14 @@ fn remove_bottleneck_in_subplan( requirements.children = requirements .children .into_iter() - .map(|node| { - if node.data { - remove_bottleneck_in_subplan(node) + .enumerate() + .map(|(idx, node)| { + // Deliberately conservative: not descending at all also skips legitimate + // cleanups *below* a protected child (a redundant second coalesce under the + // load-bearing one, say). This could later be narrowed to "descend, but + // protect only the topmost coalesce" if that turns out to matter. + if node.data && removable(idx) { + remove_bottleneck_in_subplan_impl(node, false) } else { Ok(node) } diff --git a/datafusion/sqllogictest/test_files/joins.slt b/datafusion/sqllogictest/test_files/joins.slt index e0be63fe71525..ce71d0038f747 100644 --- a/datafusion/sqllogictest/test_files/joins.slt +++ b/datafusion/sqllogictest/test_files/joins.slt @@ -5527,3 +5527,121 @@ DROP TABLE t1; statement ok DROP TABLE t2; + +# Regression test for a LEFT JOIN with a non-equijoin predicate (forces +# NestedLoopJoinExec) and a multi-partition probe side. Previously the unmatched +# left rows could be emitted before all partitions finished probing, adding +# spurious NULL-padded rows. The result must include every left row exactly once. +statement ok +set datafusion.execution.target_partitions = 4; + +statement ok +set datafusion.execution.batch_size = 2; + +statement ok +CREATE TABLE nlj_left(id INT, v INT) AS VALUES (1, 4), (2, 72), (3, 41), (4, 98), (5, 91); + +statement ok +CREATE TABLE nlj_right(w INT) AS VALUES (49), (58), (83), (3), (76); + +query III +SELECT id, v, w FROM nlj_left LEFT JOIN nlj_right ON nlj_left.v < nlj_right.w ORDER BY id, w; +---- +1 4 49 +1 4 58 +1 4 76 +1 4 83 +2 72 76 +2 72 83 +3 41 49 +3 41 58 +3 41 76 +3 41 83 +4 98 NULL +5 91 NULL + +statement ok +DROP TABLE nlj_left; + +statement ok +DROP TABLE nlj_right; + +statement ok +set datafusion.execution.target_partitions = 4; + +statement ok +reset datafusion.execution.batch_size; + +# Regression test: a `CollectLeft` `HashJoinExec` requires `SinglePartition` on its build +# (left) child, and the `CoalescePartitionsExec` that satisfies it must survive the +# sort-parallelization phase of `EnsureRequirements`. It used to be removed positionally +# (the traversal descends into the join because the *probe* side is linked to a coalesce), +# leaving a multi-partition build side that `SanityCheckPlan` rejects with +# "does not satisfy distribution requirements: SinglePartition". + +statement ok +set datafusion.execution.target_partitions = 8; + +# Keep the scan multi-partition as written, i.e. one partition per file. +statement ok +set datafusion.optimizer.repartition_file_scans = false; + +statement ok +CREATE TABLE collect_left_src (id INT, ts INT) AS VALUES (1, 10), (2, 20), (3, 30); + +query I +COPY (SELECT * FROM collect_left_src) TO 'test_files/scratch/joins/collect_left/0.parquet' STORED AS PARQUET; +---- +3 + +query I +COPY (SELECT * FROM collect_left_src) TO 'test_files/scratch/joins/collect_left/1.parquet' STORED AS PARQUET; +---- +3 + +query I +COPY (SELECT * FROM collect_left_src) TO 'test_files/scratch/joins/collect_left/2.parquet' STORED AS PARQUET; +---- +3 + +query I +COPY (SELECT * FROM collect_left_src) TO 'test_files/scratch/joins/collect_left/3.parquet' STORED AS PARQUET; +---- +3 + +statement ok +CREATE EXTERNAL TABLE collect_left STORED AS PARQUET LOCATION 'test_files/scratch/joins/collect_left/'; + +# The build side is the 4-partition scan; the probe side is the `DISTINCT ON` aggregate, +# whose `CoalescePartitionsExec` is what makes the traversal reach the join. +query I +SELECT a.id +FROM collect_left a +LEFT JOIN (SELECT DISTINCT ON (id) id, ts FROM collect_left ORDER BY id, ts) f + ON a.id = f.id +ORDER BY a.id; +---- +1 +1 +1 +1 +2 +2 +2 +2 +3 +3 +3 +3 + +statement ok +DROP TABLE collect_left; + +statement ok +DROP TABLE collect_left_src; + +statement ok +reset datafusion.optimizer.repartition_file_scans; + +statement ok +set datafusion.execution.target_partitions = 4;