Repository navigation
feat: run physical requirement enforcement first, as a PhysicalAnalyzer phase - #25688
zhuqi-lucas wants to merge 23 commits into
Conversation
3ec7e05 to
71adda4
Compare
|
Thanks @alamb — this is the first step from #25355: To keep it a pure refactor I run the analyzer phase at Would appreciate an early look at the trait/API shape. One thing I'm unsure about: a pure split looks hard in general — it's not just For the systematic version I'd build on two of your own suggestions: an |
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Critical issues remain with custom optimizer boundary placement and ForeignSession analyzer delegation.
Get a fresh assessment by requesting another Copilot review.
Review effort: Lite
Findings: 2
Open (4)
What changed in this PR
Introduces a physical analyzer phase for enforcement rules while preserving default optimizer ordering.
Changes:
- Adds
PhysicalAnalyzerRuleandPhysicalAnalyzer. - Integrates analyzers through sessions, builders, and planning.
- Moves
EnsureRequirementsto the analyzer phase. - Updates optimizer documentation and tests.
| File | Summary |
|---|---|
datafusion/session/src/session.rs |
Adds the session analyzer API. |
datafusion/session/src/physical_analyzer.rs |
Defines the analyzer rule trait. |
datafusion/session/src/lib.rs |
Exports analyzer APIs. |
datafusion/physical-optimizer/src/optimizer.rs |
Removes enforcement from the default optimizer list. |
datafusion/physical-optimizer/src/lib.rs |
Exports analyzer components. |
datafusion/physical-optimizer/src/ensure_requirements/mod.rs |
Adds analyzer compatibility for enforcement. |
datafusion/physical-optimizer/src/analyzer.rs |
Defines the default physical analyzer. |
datafusion/core/src/physical_planner.rs |
Runs analyzer and optimizer rules. |
datafusion/core/src/optimizer_rule_reference.rs |
Tests rule ordering documentation. |
datafusion/core/src/optimizer_rule_reference.md |
Documents analyzer and optimizer rules. |
datafusion/core/src/execution/session_state.rs |
Adds analyzer state and builder APIs. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
876c676 to
e618a80
Compare
I think we should consider aiming eventually for the invariant that each PhysicalOptimzierRule leaves the plan in a valid shape (so we don't have run multiple enforcement passes). That would be the simplest to reason about -- and would let each PhysicalOptimizerRule assume it got a valid plan as input which would also likely simplify the logic I am not sure how far away from this invariant the existing rules are
I think it would be even cheaper if we didn't have to re-run Enforcement |
alamb
left a comment
There was a problem hiding this comment.
This looks great @zhuqi-lucas -- in my mind, the only major thing to work out is getting all the analyzer rules to run before the optimizer rules (see comments)
I think we could also pare down some of the comments in this PR and I again left some suggestions on how to do it
| /// Responsible for optimizing a logical plan | ||
| optimizer: Optimizer, | ||
| /// Responsible for enforcing invariants on a physical execution plan | ||
| /// (distribution, ordering) before optimization |
| ret.field("query_planners", &self.inner.query_planner) | ||
| .field("analyzer", &self.inner.analyzer) | ||
| .field("optimizer", &self.inner.optimizer) | ||
| .field("physical_analyzers", &self.inner.physical_analyzers) |
There was a problem hiding this comment.
I know I am biased, but I really like the symmetry here with the analyzer
| // distribution and ordering requirements every operator declares. | ||
| // | ||
| // Ideally analyzers would run strictly first, but the built-in pipeline | ||
| // cannot: `JoinSelection` is an optimizer that decides broadcast |
There was a problem hiding this comment.
I think we should strive to get rid of this code (perhaps by refactoring the JoinSelection pass to update the distribution itself (by calling a helper of Enforce, for example) and actually put enforcement first.
I think longer term being able to easily adjust join orders and update the plan as necessary is important for having better join optimizer support
So maybe we need to start with a PR to run JoinSelection after Enforcement, and then this PR to pull enforcment to an analyzer rule that runs first
The other thing we could do temporarily is maybe put JoinSelection as an analyzer rule (as strange as that is) and file a ticket to make it a real analyzer rule
| /// Create a new analyzer using the recommended list of rules | ||
| pub fn new() -> Self { | ||
| let rules: Vec<Arc<dyn PhysicalAnalyzerRule + Send + Sync>> = vec![ | ||
| // Ensures each input plan satisfies the distribution and ordering |
There was a problem hiding this comment.
I think this comment should maybe be on PhysicalAnalyzer or the PhysicalAnalyzerRule
| // Otherwise, this rule inserts the necessary repartitioning and sorting | ||
| // operators. | ||
| // | ||
| // This used to be implemented as two separate rules: `EnforceDistribution` |
There was a problem hiding this comment.
it isn't clear to me why the history of having 2 separate rules is important for future readers to know (I think we could pare this one down)
| /// to match [`ExecutionPlan::required_input_distribution`] or insert a | ||
| /// `SortExec` to match [`ExecutionPlan::required_input_ordering`]. | ||
| /// | ||
| /// This mirrors the logical layer's split between `AnalyzerRule` (make the |
There was a problem hiding this comment.
I recommend making these doc links to ANalyzerRule and OptimizerRule (they probably need to be links to the docs.rs page due to crate dependencies)
| /// | ||
| /// Analyzer rules run as their own phase, conceptually before the optimizer | ||
| /// rules that assume a valid plan. Note that the built-in planner does not | ||
| /// literally run every analyzer before every optimizer: to preserve the |
There was a problem hiding this comment.
As mentioned elsewhere I think we should change this
| /// A flag to indicate whether the physical planner should validate that | ||
| /// the rule will not change the schema of the plan after the rewrite. | ||
| /// | ||
| /// This mirrors [`PhysicalOptimizerRule::schema_check`]; enforcement passes |
There was a problem hiding this comment.
FWIW I think schema_check should be true for all optimizer rules and only analyzer rules should be alloed to change the schema, given the definition of analyzer / optimizer we are adding. We could try and tighten this up as a follow on issue / work -- no need to do it here
Thanks @alamb, really appreciate the thorough review — this is very helpful and gives me a clear picture of the direction. I'll iterate on the PR from here (enforcement-first, plus trimming the comments) and follow up. |
Mirror the logical AnalyzerRule/OptimizerRule split on the physical side. - Add a `PhysicalAnalyzerRule` trait and `PhysicalAnalyzer` rule list, plumbed through Session / SessionState / SessionStateBuilder the same way physical optimizer rules are. - Split `EnsureRequirements` into `enforce_distribution_requirements` (Phases 0-2a), `enforce_requirements` (Phases 0-2), and `optimize_sorts` (Phase 3). It now runs as the analyzer doing *distribution enforcement only*; a new `OptimizeSorts` optimizer rule re-enforces and runs the sort optimizations exactly once, at the former `EnsureRequirements` position. The `PhysicalOptimizerRule` impl is retained (running both halves) so downstream chains that register it keep working. - The planner runs the analyzer phase first, then the optimizer phase. - `LimitedDistinctAggregation` now looks through `RepartitionExec` / `CoalescePartitionsExec`, so the distribution-first analyzer no longer drops its pushed-down limit hint. Full sqllogictest suite passes with no plan changes (only EXPLAIN VERBOSE pass-trace churn in explain.slt, regenerated). Prerequisite for the convergence-loop work in apache#25572.
|
Thanks @zhuqi-lucas -- this is going to be great |
…orceSorting / OptimizeSorts Split the monolithic enforcement into three focused rules so the phases that were conflated inside `EnsureRequirements` are explicit: - `EnforceDistribution` (Phases 0-2a): distribution enforcement only. It is the `PhysicalAnalyzerRule` that runs first to make the plan distribution-valid, and is re-registered as a `PhysicalOptimizerRule` after `JoinSelection` / `WindowTopN` change distribution. - `EnforceSorting` (Phase 2b): ordering enforcement. Not idempotent, so it runs exactly once, after the rules that settle ordering requirements. - `OptimizeSorts` (Phase 3): the sort/distribution optimizations (parallelize sorts, order-preserving variants, sort pushdown, partial sort). `EnsureRequirements` stays as a compatibility shim (distribution + sorting enforcement + sort optimization in one pass) for downstream pipelines that splice it in by position; it is no longer in either default list. `LimitedDistinctAggregation` now looks through `RepartitionExec` / `CoalescePartitionsExec` so the limit hint still reaches the partial aggregate once distribution enforcement runs first and separates it from the final aggregate. Verified zero plan churn across the sqllogictest suite. Adds tests that the three rules in sequence reproduce the monolith, that `EnforceDistribution` is idempotent on distribution-shaped plans and behaves identically as analyzer and optimizer, that it does not enforce ordering, and that the limit look-through survives a distribution operator.
e618a80 to
216b9aa
Compare
…aintaining rules
Distribution enforcement runs once, first, in the physical analyzer phase, and
every optimizer rule that changes the required distribution restores it itself
("whoever breaks it fixes it") -- so there is no second, standalone enforcement
pass and no plan churn.
- analyzer = [EnforceDistribution]: the single rule that makes the plan
distribution-valid before any optimizer rule runs. Bare root-level file scans
(e.g. `SELECT * FROM t`, which have no parent to drive the per-child split)
are parallelized by the source itself at the end of the pass.
- JoinSelection, WindowTopN and FilterPushdown each call the distribution-enforce
helper after they change the plan (only when they actually changed it):
JoinSelection tightens join distribution, WindowTopN drops the exchange under
its rewrite, and FilterPushdown flips a source's row-count stats Exact ->
Inexact by pushing a predicate, which changes parallelism decisions. Each
re-establishes distribution locally instead of relying on a later pass.
- The distribution heuristic is unchanged, so parallelism decisions match the
previous pipeline exactly.
Ordering enforcement (EnforceSorting) and the sort optimizations (OptimizeSorts)
still run once in the optimizer phase (ordering enforcement is not idempotent and
reads the settled partitioning). EnsureRequirements stays as a compatibility shim.
Result: the full sqllogictest suite is unchanged except EXPLAIN VERBOSE rule
traces in explain.slt (regenerated); every rendered plan is byte-identical, so
there is no behavior/parallelism regression. Issue apache#21826 (JoinSelection breaking
distribution after enforcement) is fixed at the source. FilterPushdown's
self-maintain updates its physical_optimizer integration snapshots (its output is
now distribution-normalized), and the join_selection tests look through the
enforcement wrappers the self-maintaining rules now add.
Pre-existing, unrelated: four memory_limit tests fail identically on the base
commit (a DiskManager/spill environment issue), not from this change.
Complete the analyzer/optimizer split so that *all* enforcement runs first, in the analyzer phase (mirroring the logical Analyzer/Optimizer split), not just distribution. Analyzer (makes the plan valid, first): - OutputRequirements(add): establish the output-requirement boundary so enforcement can see it (top-level scan parallelism, final-ordering preservation). - EnforceDistribution, then EnforceSorting on the distribution-fixed plan. Optimizer (optimization + local self-repair): - OptimizeSorts keeps only the sort optimizations, which are not enforcement. - Rules that change requirements re-establish validity themselves, scoped to exactly what they break: JoinSelection and FilterPushdown re-enforce distribution only; WindowTopN re-enforces distribution and ordering. Each repair runs only when the rule actually rewrote the plan. Golden updates: explain.slt reflects the new rule trace (enforcement now appears first, in the analyzer phase); one window_topn.slt case plans a strictly simpler, valid, correct plan (the window runs in partitioned mode over the hash-partitioned PartitionedTopKExec, eliding a redundant SortExec + SPM).
…t_partitions Two CI fixes for the analyzer-phase split. 1. OutputRequirementExec::execute delegated to its child instead of unreachable!(). Moving the OutputRequirements *add* pass into the analyzer while the *remove* pass stays in the optimizer means the marker can reach execution whenever a caller replaces the physical optimizer rules (dropping remove) but keeps the default analyzer that adds it -- e.g. the memory_limit AccessLog scenario sets the optimizer list to just [JoinSelection]. The node carries its child's partitioning/ordering, so executing the child is correct; a surviving planning marker now behaves transparently rather than panicking. Fixes the memory_limit `internal error: entered unreachable code` failures. 2. Pin target_partitions in the filter_pushdown OptimizationTest harness. Now that FilterPushdown re-establishes distribution, its snapshots can contain a RoundRobinBatch(target_partitions) repartition; with the default (num_cpus) the golden count differed between the machine snapshots were accepted on and the CI runner. Pinning it to 4 makes the goldens machine-independent; snapshots re-accepted (partition-count only).
…al plan The FFI query-planner round-trip test and the proto unoptimized-roundtrip test clear the physical optimizer rules to observe a raw plan. Now that the OutputRequirements *add* pass runs in the analyzer phase (its *remove* pass is among the optimizer rules these tests clear), the plan root is wrapped in an OutputRequirementExec marker that leaks into the FFI/proto round-trip and is rejected by the extension codec. Clear the analyzer rules too so these tests observe planning alone, matching their stated intent.
Same machine-dependent snapshot issue as the filter_pushdown harness: the WindowTopN and JoinSelection distribution self-maintain introduce a RoundRobinBatch/Hash repartition sized to target_partitions, which defaults to num_cpus. The goldens were accepted on a 12-core machine and failed on the 4-core CI runner. Pin target_partitions to 4 in both test helpers and re-accept the snapshots (partition-count only).
…ules The AccessLog / AccessLogStreaming / GroupedDistinctStrings scenarios replace the physical optimizer rules with a minimal set to keep memory-hungry repartition/sort operators out of the plan, so each test measures a specific operator's budget. Two effects of the enforcement-in-analyzer split reintroduced those operators: - Enforcement now runs in the analyzer phase, which the scenarios did not clear. Clear the analyzer rules alongside the optimizer rules so the scenario's intent (no enforcement-added repartition/sort) holds, matching the pre-split behavior. - JoinSelection now self-maintains distribution when it changes a join, adding a Hash repartition sized to target_partitions. For symmetric_hash_join that tipped the tiny budget into the spill path; pin target_partitions=1 there so no repartition is added and the test keeps measuring the join itself.
- Rewrite the PhysicalAnalyzerRule / PhysicalAnalyzer ordering docs: the analyzer phase now runs strictly before the optimizer phase, so an optimizer rule can assume a valid plan, and the rules that change a requirement (JoinSelection, WindowTopN, FilterPushdown) re-establish validity themselves. Removes the now-outdated "does not literally run every analyzer before every optimizer" caveat. - Make the logical AnalyzerRule / OptimizerRule references docs.rs doc links.
|
run benchmark sql_planner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing physical-analyzer-rule (d5ce43a) to 9ee892b (merge-base) diff Run configurationrun benchmark sql_plannerResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing physical-analyzer-rule (d5ce43a) to 9ee892b (merge-base) diff Run configurationrun benchmark sql_plannerCPU Details (lscpu)Details
Resource Usagesql_planner — base (merge-base)
sql_planner — branch
File an issue against this benchmark runner |
So one way to bring back performance then is to update JoinSelection so it doesn't have to call Sorry I find it really hard to read large / wall of comments, so you may have already said this farther down
I guess what I am advocating is trying to introduce PhysicalAnalyzerRule with the smallest number of other changes as possible. For example, if we had to initially treat Then in follow on PRs we could explore how to move JoinSelection into the OptimizerRules (e.g. by ensuring that it doesn't make plans invalid) Those steps I think are more self contained and easier to review |
The module diagram still showed Phase 2 as one combined distribution and sorting pass; the code has always been a distribution walk (2a) followed by an ordering walk (2b), and the rest of the module now names them that way.
… order Per review: introduce PhysicalAnalyzerRule with the smallest possible change. The default analyzer list is the former optimizer list up to and including EnsureRequirements, in the same order, so no plan changes, not even the EXPLAIN VERBOSE trace. Every rule keeps its logic; the seven analyzer rules delegate to their existing PhysicalOptimizerRule impls. Dropped from this PR, to be proposed separately: the EnsureRequirements decomposition and OptimizeSorts, moving AggregateStatistics / LimitedDistinctAggregation / FilterPushdown / WindowTopN after enforcement, the JoinSelection phase enum, and the debug-only boundary invariant check.
|
Thanks @alamb, and sorry for the long comments, I will keep them short from now on.
Done, the PR is now exactly that: the trait, the plumbing, and the default analyzer list is today's optimizer list up to and including And the benchmark is good, no regression: The follow-ups are listed in the PR description: moving |
eb879ec to
86fde11
Compare
|
run benchmark sql_planner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing physical-analyzer-rule (79c63da) to 4d167a1 (merge-base) diff Run configurationrun benchmark sql_plannerResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing physical-analyzer-rule (79c63da) to 4d167a1 (merge-base) diff Run configurationrun benchmark sql_plannerCPU Details (lscpu)Details
Resource Usagesql_planner — base (merge-base)
sql_planner — branch
File an issue against this benchmark runner |
|
Here is what I understand from the latest discussion and code diff. I used AI to help summarize it, so please point out anything I misunderstood. Summary of current implementationThe latest change keeps the existing physical optimizer rule order, but introduces one hard boundary through the Conceptually:
My main suggestion: Extensible BoundariesIt would be very helpful to make this boundary mechanism extensible, instead of encoding one specific boundary through the Analyzer/Optimizer split. If these boundaries can be represented explicitly, I think optimizer maintenance becomes much easier. Rules can clearly see which assumptions they are allowed to rely on, instead of those assumptions being hidden in rule ordering and implementation details. Projection pushdown is one concrete example of where another boundary could be useful. The important part is not this specific invariant, but that the same mechanism could express it naturally: I suspect there are more hidden invariants like this in the optimizer today, so I think it is important that the mechanism can naturally support multiple boundaries rather than only this one distribution/ordering boundary. Nice-to-have: explicit validations for boundariesI don't think this needs to be a major design decision, since it seems relatively easy to extend once boundaries are represented explicitly. For example, after a boundary establishes an invariant, we could validate it after every subsequent rule, at least in tests or debug builds: That would make it much easier to identify exactly which rule first breaks the invariant, rather than discovering the violation only after the entire optimizer finishes. |
|
Thanks @2010YOUY01, your summary is accurate. I like the extensible boundaries idea. The analyzer/optimizer split is the first such boundary, and a generic "boundary + validator" entry in the rule list can build on it rather than replace it. I would treat this PR as the starting point and do the general mechanism as a follow-up, there is a lot we can build on top of it. The validation part is cheap: behind |
Thank you @2010YOUY01 and @zhuqi-lucas Extensible boundaries are an interesting idea. The one usecase I know of is a OptimizerRule that may disturb sorting or distribution requirements (and thus needs to run the EnforceRequirements pass again). I think this usecase can be satisfied by having the OptimizerRule just call EnforceRequirements again explicitly so the requirements are satisfied after the rule runs. As neat an idea of having multiple classes of constraints that can be satisfied at different points, I think it makes it harder to reason about the optimizer as a whole and any OptimizerRule individually -- we would need some way for each rule to communicate what its expected input and output boundary was I also don't see the split of "invalid plan" --> "valid plan" as embodied in AnalayzeRule / OptimizerRule as preventing us from adding some sort of extensible boundaries in the future (though they would likely not have an invalid/valid boundary) 🤔 |
I don't fully understand the use of this boundary -- when writing an OptimizerRule or trying to understand one, ideally I would not care if the projection was an explicit ProjectionExec or embedded in another ExecutionPlan unless it was relevant for my optimization |
The short answer is: if we can define a region in the optimizer rule list where projections are guaranteed to have a canonical shape, then every rule within that region only needs to pattern-match against that single shape, instead of handling all possible projection variants. This can simplify the implementations. There are many similar regions that already exist implicitly today, so the simplification could add up substantially. Existing implementations are often unaware of these implicit regions, so they are written more defensively and end up carrying extra complexity. Here is one example: #25688 (comment) (Below is a more detailed explanation of my thinking, which currently leans toward implementing this optimizer boundary mechanism directly.)
If I understand the assumptions behind the optimizer implementation correctly, adding more explicit constraints should make both understanding and implementation easier. Let me first restate the assumptions. I may be missing something here since I have less experience working on optimizers. # Assumptions for the physical optimizer
DataFusion provides a default ordered list of optimizer rules.
// Default rules
rules = [
rule1,
rule2,
rule3,
rule4,
...
]
- The default rules are only guaranteed to be safe when run in the prescribed order.
- To disable specific behavior, use configuration options that define known-safe variations. For example,
`set datafusion.optimizer.enable_distinct_aggregation_soft_limit = false`
disables one specific aggregation optimization without arbitrarily modifying the rule pipeline.
- (This sounds intimidating, but seems close to the current state):
To implement an extension rule outside core, you need to understand the implicit constraints established by the surrounding default optimizer rules. Otherwise, changes in the default rules might break the plan.
Currently, many of these constraints are only documented locally inside individual rules.
# Future improvement direction
Make these constraints easier to understand, express, and test.
# Not intended usage
Arbitrarily reorder/remove default rules downstream and expect the resulting optimizer pipeline to remain valid.
// Reordered default rules mixed with extension rules
rules = [
rule3,
extension_rule1,
rule1,
...
]More constraints can simplify optimizer extensionsIf we explicitly exclude arbitrary reordering of the default optimizer rule list from the supported use case, then stronger constraints can simplify the optimizer:
WalkthroughIn @zhuqi-lucas's use case: rules = [
rule1,
rule2,
EnsureRequirements,
rule3,
extension_rule,
...
]If we treat the analyzer/optimizer split as a general boundary, the preconditions and postconditions of
The simplification provided by this boundary is that a semantically equivalent plan has a canonical physical shape at this point in the optimizer. // plan1 and plan2 are semantically equivalent, but the analyzer boundary
// guarantees that this rule only needs to handle plan1.
//
// This reduces the number of physical shapes the rule needs to reason about.
// plan1
AggregateExec(mode=final)
--RepartitionExec
----AggregateExec(mode=partial)
// plan2
AggregateExec(mode=final)
--AggregateExec(mode=partial)More constraints can reduce the problem space further. For example, if this extension rule needs to inspect projections, its implementation becomes simpler if projected plans also have one canonical shape at this boundary, rather than allowing several equivalent representations. Proceeding with this PRNow @zhuqi-lucas and @alamb seem more in favor of implementing the analyzer/optimizer split first, and treating a more general boundary mechanism as a follow-up. My current preference is to implement the more general optimizer boundary directly, to avoid introducing multiple mechanisms for the same underlying problem. That said, if I have misunderstood the assumptions above, this proposal would definitely need to be reworked. I'd like to hear your thoughts before deciding on the next step. If you also think a general boundary mechanism is a plausible direction, I can put together a PoC this week. |
There happens to be an example of how an extra constraint can simplify the implementation; see the rationale for details. The TL;DR is that the existing optimizer rewrite implementation tries to recognize 3 possible aggregate plan shapes, and does so using a nested DFS algorithm. However, there is a hidden constraint at the point where this rule runs, which means only one of those shapes is actually possible. With that constraint made explicit, the rewrite can be simplified to a straightforward pattern match without nested DFS. |
Thank you -- this is a very helpful example. I suspect the reason the rule look for three possible shapes is that they used to be created in earlier versions of DataFusion, so I wonder what type of boundary we could use to encode this in the future.
Yes, this accurately describes my preference, though I am happy to explore other alternatives and don't see any reason to rush this in. While the idea of a general boundary mechanism seems good in theory, I am somewhat skeptical of how useful it would be in practice. Specifically, I struggle to imagine what specific boundaries / invariants we could encode in an enumeration an enforce during planing other than the existing ones:
For example it is not clear to me the property described after I think my main concerns are:
If we can come up with some other specific, useful boundaries / invariants this sounds like a good plan to me. |
In my mind, the constraints are not imposed by the rules themselves, but are a property of the ExecutionPlan nodes -- there should be clear semantics of what each ExecutionPlan does and what is required of its input/output As long as these constraints are satisfied, then an OptimizerRule should be able to rearrange the plan however it sees fit as long as the resulting plan still produces the same output What is not clear to me is that there is any way to encode some sort of optimization property (e.g that limits have been pushed down) that gets tighter after each pass |
|
Thank you @alamb and @2010YOUY01. One observation from reading #26069 against the code: its "only shape 1 is possible" precondition holds exactly because On the general mechanism, I agree with @alamb's framing: the invariants a plan can check about itself (schema, declared input properties satisfied) are what this boundary encodes, and I would land this PR as the base, and if a second concrete boundary that the plan itself can check shows up, build the general mechanism on top of it. Happy to review the PoC if @2010YOUY01 wants to explore it in parallel. |
|
My current intuition is that we can keep refactoring the existing optimizer by making more invariants and boundaries explicit, as in #26069, and use them to continuously reduce overall complexity. I think the concerns from @alamb and @zhuqi-lucas are:
Those concerns make a lot of sense to me. My next step is to prototype the idea and try to identify/enforce ~3 more useful boundaries or invariants, to see whether this approach can consistently reduce complexity:
|
Thank you for the measured and thoughtful reply. I think the idea of prototyping out what it might look like is a very reasonable next step |


Which issue does this PR close?
Rationale for this change
The physical optimizer runs a single hand-ordered rule list. Some of those passes are not optimizations, they are enforcers of invariants: they insert the repartitioning and sorting needed to satisfy the distribution and ordering requirements every operator declares. The logical layer already separates these two kinds of passes (
AnalyzerRulemakes a plan valid,OptimizerRulemakes it faster); the physical layer does not.This PR introduces that split with the smallest possible change: a
PhysicalAnalyzerRuletrait and aPhysicalAnalyzerphase that runs before thePhysicalOptimizer. The default analyzer list is today's optimizer list up to and includingEnsureRequirements, in the same order, so no plan changes. EveryPhysicalOptimizerRulenow receives a plan whose requirements are already enforced.Architecture
Rules 2 to 6 sit in the analyzer only because enforcement has to see their output today (
join_selectionresolvesPartitionMode::Auto, which declares no requirement; the others rewrite aggregates and windows that enforcement then repartitions and sorts). Moving them into the optimizer phase, one at a time, is follow-up work (see below).What changes are included in this PR?
PhysicalAnalyzerRuletrait and aPhysicalAnalyzerrule list, plumbed throughSession/SessionState/SessionStateBuildersymmetrically with the physical optimizer rules. The planner runs the analyzer phase, then the optimizer phase.impl PhysicalAnalyzerRulefor the seven rules in the default analyzer list; each delegates to its existingPhysicalOptimizerRuleimpl, so no rule logic changes.OutputRequirementExec::executedelegates to its input instead of panicking, since its remove pass now lives in a different phase than its add pass and a custom optimizer list can drop it.optimizer_rule_reference.mdgains the pipeline diagram and aPhysical Analyzer Rulessection; the upgrade guide documents the new phase.Follow-ups
Each of these is a self-contained PR on top of this one:
aggregate_statistics,LimitedDistinctAggregationandFilterPushdowninto the optimizer phase.aggregate_statisticsalready looks through exchanges;LimitedDistinctAggregationneeds to look throughRepartitionExec/CoalescePartitionsExec;FilterPushdownonly changes the parallelism decision for sources under 8192 rows.WindowTopNinto the optimizer phase: it changes the window's input requirements, so it has to re-establish them for the subtree it rewrote; the sort optimizations inEnsureRequirements(phase 3) then need to run after it as their own rule.join_selection: resolvingPartitionMode::Autois what enforcement depends on, so either enforcement resolves it inline orJoinSelectionrepairs distribution locally after enforcement. Needs its own design discussion.InvariantChecker(InvariantLevel::Executable)in debug builds, and teachHashJoinExec::check_invariantsto rejectPartitionMode::Auto.Are these changes tested?
Yes:
PhysicalAnalyzerRuleregistered via the builder runs during planning.optimizer_rule_reference.mdis updated and its parsing tests pin the documented analyzer/optimizer order to the code.sqllogictestsuite passes with no golden changes: the rules run in the same order as before, so even theEXPLAIN VERBOSEtrace is unchanged.Are there any user-facing changes?
PhysicalAnalyzerRule,PhysicalAnalyzer,Session::physical_analyzers,SessionState::physical_analyzers/add_physical_analyzer_rule, andSessionStateBuilder::with_physical_analyzer_rule(s).DefaultPhysicalPlanner's optimize observer callback now receives the rule name (&str) instead of&dyn PhysicalOptimizerRule, so analyzer passes are surfaced too (for example inEXPLAIN VERBOSE).Sessionthat overridesphysical_optimizers()should also overridephysical_analyzers(); a custom optimizer list that still contains a rule now in the analyzer runs it twice unless the analyzer list is cleared.