adriangb opened a new pull request, #23948:
URL: https://github.com/apache/datafusion/pull/23948

   ## Which issue does this PR close?
   
   - None filed; happy to open one if preferred.
   
   ## Rationale for this change
   
   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.
   
   Reproducer with `datafusion-cli` (write four small parquet files into 
`data/` first, so the scan is multi-partition):
   
   ```sql
   set datafusion.execution.target_partitions = 8;
   set datafusion.optimizer.repartition_file_scans = false;
   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.
   
   ## What changes are included in this PR?
   
   `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.
   
   ## Are these changes tested?
   
   Yes: two new tests in 
`datafusion/core/tests/physical_optimizer/ensure_requirements.rs` cover 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. Both 
fail on `main` with the distribution error above.
   
   `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.
   
   ## Are there any user-facing changes?
   
   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.
   


-- 
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]

Reply via email to