zhuqi-lucas commented on PR #25688: URL: https://github.com/apache/datafusion/pull/25688#issuecomment-5932417602
> I think the structure of this PR looks good in general. I worry about: > > 1. Plan changes -- they don't all look good to me > 2. Performance regression (can we perhaps combine EnforceDistribution and EnforceSort again and find some other ways to make OptimizeSort faster?) Thanks @alamb. On 2, I profiled it, and I do not think combining them back is the lever. Local A/B on tpch q21, merge-base `e1aa7d9` against this branch, same machine: | | merge-base | this branch | |---|---|---| | end to end per plan | 1173.2 us | 1383.2 us | | `join_selection` | 80.7 us | **252.9 us** | | `EnsureRequirements` | 133.4 us | moved to the analyzer | | `OptimizeSorts` | - | 29.8 us | The optimizer-phase delta is `join_selection` +172.2, `EnsureRequirements` -133.4, `OptimizeSorts` +29.8. So the extra time is inside `JoinSelection`, and it is the `enforce_distribution_requirements` call it makes after rewriting. That function runs **twice per plan** here and is 22% of planning. Splitting the two rules is not where it went: on merge-base `EnsureRequirements` already does distribution and sorting as two separate bottom-up walks, so merging them removes no traversal, and `OptimizeSorts` at 29.8 us cannot account for a 210 us regression even at zero. Looking at the three rules that self-repair in this PR, they are not the same case: | rule | repairs | why it touches enforcement | where it belongs | |---|---|---|---| | `JoinSelection` | distribution | declares nothing until it runs | **before** enforcement | | `FilterPushdown` | distribution | see below, two reasons | **before** enforcement | | `WindowTopN` | distribution and ordering | needs the lowered plan to do its job | **after**, repair is legitimate | `JoinSelection`: before it runs, a hash join's mode is `PartitionMode::Auto`, which declares `UnspecifiedDistribution` for its children and `UnknownPartitioning` for its output. Enforcement running first therefore does nothing for joins. The plan only acquires unsatisfied requirements once a mode is picked, and the repair then re-derives them over the whole tree, on a tree the first pass already expanded. `FilterPushdown` is coupled to enforcement in two separate ways, and I think only one of them is real. First, pushing a predicate into a source rebuilds `DataSourceExec`, and the rebuilt node loses the file grouping enforcement had established: in `repartition_scan.slt` 4 file groups become 2 when I remove the repair. That looks like a defect in the rebuild rather than an ordering constraint, and fixing it would decouple the two. Second, pushing a predicate changes the statistics the distribution decisions read, a source's row count flips Exact to Inexact, which changes whether its scan is parallelized. That one is a genuine dependency and is why master runs `FilterPushdown` before `EnsureRequirements`. `WindowTopN` is the one that earns the contract: it rebuilds the window over a hash-partitioned `PartitionedTopKExec`, so it needs enforcement to have run, and it genuinely disturbs both distribution and ordering. It is also rare enough that it does not show up in the q21 numbers at all. I tried moving `JoinSelection` ahead of enforcement rather than repairing after it. q21 goes from **1.3973 ms to 1.1856 ms** against the merge-base 1.2312 ms, so faster than before this PR rather than merely recovered. The `issue_21826` regression test still passes with only the `InterleaveExec` normalisation kept inside `JoinSelection`, which is the part that actually fixes that panic and holds in any order. To be explicit, that makes `JoinSelection` an analyzer rule rather than an optimizer one. It is what you floated earlier in this review, "put `JoinSelection` as an analyzer rule (as strange as that is)", and I think it is less strange than it sounds: `PartitionMode::Auto` is not executable at all, `HashJoinExec::execute` returns a plan error on it, so resolving the mode is a correctness step under the definition you gave. Where I am less sure, and where @2010YOUY01's staging seems relevant: if both `JoinSelection` and `FilterPushdown` have to precede enforcement while `WindowTopN`, `OptimizeSorts` and `TopKRepartition` have to follow it, then a plain analyzer-then-optimizer split does not express the chain, since `FilterPushdown` is an optimization by any reading. Either the analyzer phase ends up holding rules that are not about correctness, or the split wants a different criterion than correctness, something closer to "determines requirements" versus "consumes them". I would rather measure the first part and let the naming follow than argue it up front, but I would like to know how you both read that boundary before I rework this. On 1, I replied inline about the `window_topn` plan, where parallelism goes up rather than down. Happy to walk through the rest of the plan changes one at a time. -- 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]
