Skip to content

feat: run physical requirement enforcement first, as a PhysicalAnalyzer phase - #25688

Open
zhuqi-lucas wants to merge 23 commits into
apache:mainfrom
zhuqi-lucas:physical-analyzer-rule
Open

zhuqi-lucas wants to merge 23 commits into
apache:mainfrom
zhuqi-lucas:physical-analyzer-rule

Conversation

@zhuqi-lucas

@zhuqi-lucas zhuqi-lucas commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

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 (AnalyzerRule makes a plan valid, OptimizerRule makes it faster); the physical layer does not.

This PR introduces that split with the smallest possible change: a PhysicalAnalyzerRule trait and a PhysicalAnalyzer phase that runs before the PhysicalOptimizer. The default analyzer list is today's optimizer list up to and including EnsureRequirements, in the same order, so no plan changes. Every PhysicalOptimizerRule now receives a plan whose requirements are already enforced.

Architecture

ExecutionPlan  (every operator declares its distribution and ordering
                requirements; nothing satisfies them yet)
                               │
                               ▼
┌─ PhysicalAnalyzer phase ── make the plan VALID ──────────────┐
│ 1. OutputRequirements (add)   establish the output boundary  │
│ 2. aggregate_statistics       answer aggregates from stats   │
│ 3. join_selection             resolve PartitionMode::Auto    │
│ 4. LimitedDistinctAggregation push limit hints into aggs     │
│ 5. FilterPushdown             push predicates into sources   │
│ 6. WindowTopN                 window + filter into TopK      │
│ 7. EnsureRequirements         enforce distribution, ordering │
└──────────────────────────────┬───────────────────────────────┘
                               │  invariant: the plan is valid
                               ▼
┌─ PhysicalOptimizer phase ── make the plan FASTER ────────────┐
│ every rule receives a valid plan and leaves a valid plan     │
└──────────────────────────────┬───────────────────────────────┘
                               ▼
                  valid, optimized ExecutionPlan

Rules 2 to 6 sit in the analyzer only because enforcement has to see their output today (join_selection resolves PartitionMode::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?

  • A new PhysicalAnalyzerRule trait and a PhysicalAnalyzer rule list, plumbed through Session / SessionState / SessionStateBuilder symmetrically with the physical optimizer rules. The planner runs the analyzer phase, then the optimizer phase.
  • impl PhysicalAnalyzerRule for the seven rules in the default analyzer list; each delegates to its existing PhysicalOptimizerRule impl, so no rule logic changes.
  • OutputRequirementExec::execute delegates 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.md gains the pipeline diagram and a Physical Analyzer Rules section; the upgrade guide documents the new phase.

Follow-ups

Each of these is a self-contained PR on top of this one:

  1. Move aggregate_statistics, LimitedDistinctAggregation and FilterPushdown into the optimizer phase. aggregate_statistics already looks through exchanges; LimitedDistinctAggregation needs to look through RepartitionExec / CoalescePartitionsExec; FilterPushdown only changes the parallelism decision for sources under 8192 rows.
  2. Move WindowTopN into 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 in EnsureRequirements (phase 3) then need to run after it as their own rule.
  3. join_selection: resolving PartitionMode::Auto is what enforcement depends on, so either enforcement resolves it inline or JoinSelection repairs distribution locally after enforcement. Needs its own design discussion.
  4. Check the phase boundary with InvariantChecker(InvariantLevel::Executable) in debug builds, and teach HashJoinExec::check_invariants to reject PartitionMode::Auto.

Are these changes tested?

Yes:

  • Unit tests pin the default analyzer order and that none of its rules is run again by the default optimizer.
  • A planner test that a custom PhysicalAnalyzerRule registered via the builder runs during planning.
  • optimizer_rule_reference.md is updated and its parsing tests pin the documented analyzer/optimizer order to the code.
  • The full sqllogictest suite passes with no golden changes: the rules run in the same order as before, so even the EXPLAIN VERBOSE trace is unchanged.

Are there any user-facing changes?

  • New public API: PhysicalAnalyzerRule, PhysicalAnalyzer, Session::physical_analyzers, SessionState::physical_analyzers / add_physical_analyzer_rule, and SessionStateBuilder::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 in EXPLAIN VERBOSE).
  • Migration notes are in the 56.0.0 upgrade guide: a custom Session that overrides physical_optimizers() should also override physical_analyzers(); a custom optimizer list that still contains a rule now in the analyzer runs it twice unless the analyzer list is cleared.

@github-actions github-actions Bot added optimizer Optimizer rules core Core DataFusion crate labels Sep 24, 2026
@zhuqi-lucas
zhuqi-lucas force-pushed the physical-analyzer-rule branch from 3ec7e05 to 71adda4 Compare September 24, 2026 08:12
@zhuqi-lucas

zhuqi-lucas commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor Author

Thanks @alamb — this is the first step from #25355: EnsureRequirements is pulled out into a new PhysicalAnalyzerRule trait, and the planner runs analyzers then optimizers.

To keep it a pure refactor I run the analyzer phase at EnsureRequirements' current position (right before CombinePartialFinalAggregate) rather than strictly first: running it before JoinSelection yields an invalid single-partition broadcast join, so this way every plan and the full sqllogictest suite stay unchanged.

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 JoinSelection; several rules depend on where enforcement runs, some needing to be before it and others after it (e.g. CombinePartialFinalAggregate after distribution enforcement), so it feels like epic-sized work.

For the systematic version I'd build on two of your own suggestions: an enforce -> optimize -> enforce alternation (enforcement still runs on both sides, since rules depend on it in each direction), with each enforcement pass gated on check_invariants so one whose requirements already hold is skipped by a cheap semantic check instead of a full re-derivation — repeated enforcement stays correct but cheap. Happy to capture all of this (topology, the invariant gate, and the deprecated shim) in an epic so this PR stays a small prerequisite — I can open it if you agree.

@zhuqi-lucas
zhuqi-lucas marked this pull request as ready for review September 24, 2026 08:24
Copilot AI lite review requested due to automatic review settings September 24, 2026 08:24

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 High severity · 2 Low severity

Open (4)
What changed in this PR

Introduces a physical analyzer phase for enforcement rules while preserving default optimizer ordering.

Changes:

  • Adds PhysicalAnalyzerRule and PhysicalAnalyzer.
  • Integrates analyzers through sessions, builders, and planning.
  • Moves EnsureRequirements to 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.

Comment thread datafusion/core/src/physical_planner.rs Outdated
Comment thread datafusion/session/src/session.rs
Comment thread datafusion/physical-optimizer/src/analyzer.rs Outdated
Comment thread datafusion/session/src/physical_analyzer.rs Outdated
@alamb

alamb commented Sep 24, 2026

Copy link
Copy Markdown
Contributor

an enforce -> optimize -> enforce alternation (enforcement still runs on both sides, since rules depend on it in each direction), with each enforcement pass gated on check_invariants so one whose requirements already hold is skipped by a cheap semantic check instead of a full re-derivation — repeated enforcement stays correct but cheap. Happy to capture all of this (topology, the invariant gate, and the deprecated shim) in an epic so this PR stays a small prerequisite — I can open it if you agree.

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

repeated enforcement stays correct but cheap

I think it would be even cheaper if we didn't have to re-run Enforcement

@alamb alamb left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

❤️

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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I know I am biased, but I really like the symmetry here with the analyzer

Comment thread datafusion/core/src/physical_planner.rs Outdated
// 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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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`

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

@zhuqi-lucas

zhuqi-lucas commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor Author

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

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.
@alamb

alamb commented Sep 24, 2026

Copy link
Copy Markdown
Contributor

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.
@zhuqi-lucas
zhuqi-lucas force-pushed the physical-analyzer-rule branch from e618a80 to 216b9aa Compare September 25, 2026 14:23
@zhuqi-lucas
zhuqi-lucas marked this pull request as draft September 25, 2026 14:23
@github-actions github-actions Bot added the sqllogictest SQL Logic Tests (.slt) label Sep 25, 2026
…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).
@zhuqi-lucas zhuqi-lucas changed the title feat: extract enforcement into a PhysicalAnalyzerRule phase (prereq for #25572) feat: run physical requirement enforcement first, as a PhysicalAnalyzer phase Sep 26, 2026
…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.
@github-actions github-actions Bot added proto Related to proto crate ffi Changes to the ffi crate labels Sep 27, 2026
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.
@github-actions github-actions Bot added the physical-plan Changes to the physical-plan crate label Oct 1, 2026
@zhuqi-lucas

Copy link
Copy Markdown
Contributor Author

run benchmark sql_planner

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark running (GKE) | trigger
Instance: c4a-highcpu-48 (12 vCPU / 65 GiB) | Linux bench-c5934038705-3004-qs59z 6.12.94+ #1 SMP Wed Aug 19 07:47:20 UTC 2026 aarch64 GNU/Linux

CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  48
On-line CPU(s) list:                     0-47
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     48
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               3 MiB (48 instances)
L1i cache:                               3 MiB (48 instances)
L2 cache:                                96 MiB (48 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-47
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected

Comparing physical-analyzer-rule (d5ce43a) to 9ee892b (merge-base) diff

Run configuration
run benchmark sql_planner

Results will be posted here when complete


File an issue against this benchmark runner

@github-actions github-actions Bot added the auto detected api change Auto detected API change label Oct 1, 2026
@adriangbot

Copy link
Copy Markdown

🤖 Benchmark completed (GKE) | trigger

Instance: c4a-highcpu-48 (12 vCPU / 65 GiB)

Comparing physical-analyzer-rule (d5ce43a) to 9ee892b (merge-base) diff

Run configuration
run benchmark sql_planner
CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  48
On-line CPU(s) list:                     0-47
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     48
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               3 MiB (48 instances)
L1i cache:                               3 MiB (48 instances)
L2 cache:                                96 MiB (48 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-47
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected
Details

group                                                 HEAD                                   physical-analyzer-rule
-----                                                 ----                                   ----------------------
logical_aggregate_with_join                           1.00    410.3±1.65µs        ? ?/sec    1.00    411.3±1.74µs        ? ?/sec
logical_correlated_subquery_exists                    1.00    246.2±0.99µs        ? ?/sec    1.01    248.7±1.07µs        ? ?/sec
logical_correlated_subquery_in                        1.00    248.2±1.33µs        ? ?/sec    1.01    250.6±1.18µs        ? ?/sec
logical_distinct_many_columns                         1.00    330.9±1.14µs        ? ?/sec    1.00    329.3±1.06µs        ? ?/sec
logical_join_4_with_agg_and_filter                    1.00    200.1±0.64µs        ? ?/sec    1.00    200.1±0.96µs        ? ?/sec
logical_join_8_with_agg_sort_limit                    1.01    370.7±1.75µs        ? ?/sec    1.00    365.8±1.08µs        ? ?/sec
logical_join_chain_16                                 1.02    639.5±1.98µs        ? ?/sec    1.00    629.3±1.82µs        ? ?/sec
logical_join_chain_4                                  1.00     78.1±0.40µs        ? ?/sec    1.01     79.1±0.23µs        ? ?/sec
logical_join_chain_8                                  1.00    207.5±0.55µs        ? ?/sec    1.00    206.7±0.47µs        ? ?/sec
logical_multiple_subqueries                           1.00    483.1±1.39µs        ? ?/sec    1.00    483.2±1.70µs        ? ?/sec
logical_nested_cte_4_levels                           1.00    223.0±1.37µs        ? ?/sec    1.00    222.9±1.03µs        ? ?/sec
logical_plan_struct_join_agg_sort                     1.04    134.1±0.90µs        ? ?/sec    1.00    129.3±0.87µs        ? ?/sec
logical_plan_tpcds_all                                1.00     84.0±0.13ms        ? ?/sec    1.01     85.1±0.21ms        ? ?/sec
logical_plan_tpch_all                                 1.00      5.6±0.02ms        ? ?/sec    1.00      5.6±0.02ms        ? ?/sec
logical_scalar_subquery                               1.00    267.4±1.01µs        ? ?/sec    1.00    268.5±0.88µs        ? ?/sec
logical_select_all_from_1000                          1.00      8.8±0.03ms        ? ?/sec    1.00      8.8±0.02ms        ? ?/sec
logical_select_one_from_700                           1.00    295.5±1.55µs        ? ?/sec    1.00    294.1±1.60µs        ? ?/sec
logical_trivial_join_high_numbered_columns            1.00    249.4±1.14µs        ? ?/sec    1.01    252.1±1.54µs        ? ?/sec
logical_trivial_join_low_numbered_columns             1.00    236.2±1.41µs        ? ?/sec    1.00    237.1±1.31µs        ? ?/sec
logical_union_4_branches                              1.00    390.9±1.68µs        ? ?/sec    1.00    391.0±1.42µs        ? ?/sec
logical_union_8_branches                              1.00    786.3±1.92µs        ? ?/sec    1.00    784.9±1.97µs        ? ?/sec
logical_wide_aggregate_100_exprs                      1.00      3.9±0.01ms        ? ?/sec    1.00      3.9±0.01ms        ? ?/sec
logical_wide_case_50_exprs                            1.00   1799.2±6.45µs        ? ?/sec    1.01   1809.7±6.86µs        ? ?/sec
logical_wide_filter_200_predicates                    1.00   1306.0±9.76µs        ? ?/sec    1.00   1312.5±6.31µs        ? ?/sec
logical_wide_filter_50_predicates                     1.00    361.2±2.17µs        ? ?/sec    1.00    359.8±1.39µs        ? ?/sec
optimizer_correlated_exists                           1.01    227.3±0.60µs        ? ?/sec    1.00    226.1±0.49µs        ? ?/sec
optimizer_join_4_with_agg_filter                      1.03    420.3±1.63µs        ? ?/sec    1.00    409.1±1.40µs        ? ?/sec
optimizer_join_chain_4                                1.02    161.6±0.33µs        ? ?/sec    1.00    158.9±0.20µs        ? ?/sec
optimizer_join_chain_8                                1.01    548.3±1.35µs        ? ?/sec    1.00    545.5±0.81µs        ? ?/sec
optimizer_select_all_from_1000                        1.00      6.7±0.01ms        ? ?/sec    1.00      6.7±0.02ms        ? ?/sec
optimizer_select_one_from_700                         1.00    239.0±0.41µs        ? ?/sec    1.00    238.9±0.42µs        ? ?/sec
optimizer_tpcds_all                                   1.00    287.5±0.42ms        ? ?/sec    1.00    288.5±1.01ms        ? ?/sec
optimizer_tpch_all                                    1.00     15.9±0.04ms        ? ?/sec    1.00     15.9±0.13ms        ? ?/sec
optimizer_wide_aggregate_100                          1.01   1793.9±4.57µs        ? ?/sec    1.00   1783.4±3.98µs        ? ?/sec
optimizer_wide_filter_200                             1.00      3.9±0.01ms        ? ?/sec    1.01      3.9±0.01ms        ? ?/sec
physical_intersection                                 1.00    575.3±1.46µs        ? ?/sec    1.01    578.9±2.15µs        ? ?/sec
physical_join_consider_sort                           1.00   1027.4±2.26µs        ? ?/sec    1.01   1039.9±1.90µs        ? ?/sec
physical_join_distinct                                1.00    229.9±1.36µs        ? ?/sec    1.00    229.1±1.33µs        ? ?/sec
physical_many_self_joins                              1.00      7.8±0.02ms        ? ?/sec    1.00      7.8±0.02ms        ? ?/sec
physical_plan_clickbench_all                          1.00    137.6±0.25ms        ? ?/sec    1.02    139.7±0.92ms        ? ?/sec
physical_plan_clickbench_q1                           1.00   1343.2±8.35µs        ? ?/sec    1.07   1441.6±8.22µs        ? ?/sec
physical_plan_clickbench_q10                          1.00   1943.8±4.00µs        ? ?/sec    1.00   1942.8±7.01µs        ? ?/sec
physical_plan_clickbench_q11                          1.00      2.1±0.00ms        ? ?/sec    1.00      2.1±0.00ms        ? ?/sec
physical_plan_clickbench_q12                          1.00      2.2±0.00ms        ? ?/sec    1.00      2.2±0.00ms        ? ?/sec
physical_plan_clickbench_q13                          1.00   1968.6±6.42µs        ? ?/sec    1.00   1972.2±4.73µs        ? ?/sec
physical_plan_clickbench_q14                          1.00      2.1±0.00ms        ? ?/sec    1.00      2.1±0.01ms        ? ?/sec
physical_plan_clickbench_q15                          1.00      2.0±0.00ms        ? ?/sec    1.00      2.0±0.00ms        ? ?/sec
physical_plan_clickbench_q16                          1.00   1733.1±3.84µs        ? ?/sec    1.00   1729.6±4.88µs        ? ?/sec
physical_plan_clickbench_q17                          1.00   1779.9±4.27µs        ? ?/sec    1.00   1782.5±6.06µs        ? ?/sec
physical_plan_clickbench_q18                          1.00   1632.6±3.43µs        ? ?/sec    1.00   1634.6±4.83µs        ? ?/sec
physical_plan_clickbench_q19                          1.00   1977.2±4.19µs        ? ?/sec    1.00   1973.8±5.62µs        ? ?/sec
physical_plan_clickbench_q2                           1.00   1732.0±3.82µs        ? ?/sec    1.00   1735.7±5.27µs        ? ?/sec
physical_plan_clickbench_q20                          1.00   1510.8±4.58µs        ? ?/sec    1.00   1508.9±5.70µs        ? ?/sec
physical_plan_clickbench_q21                          1.00   1723.0±6.17µs        ? ?/sec    1.00   1723.4±5.38µs        ? ?/sec
physical_plan_clickbench_q22                          1.00      2.1±0.00ms        ? ?/sec    1.00      2.1±0.00ms        ? ?/sec
physical_plan_clickbench_q23                          1.01      2.2±0.00ms        ? ?/sec    1.00      2.2±0.00ms        ? ?/sec
physical_plan_clickbench_q24                          1.00      5.1±0.01ms        ? ?/sec    1.00      5.1±0.01ms        ? ?/sec
physical_plan_clickbench_q25                          1.00   1822.3±4.01µs        ? ?/sec    1.00   1822.6±4.88µs        ? ?/sec
physical_plan_clickbench_q26                          1.00   1687.8±4.44µs        ? ?/sec    1.00   1692.0±5.66µs        ? ?/sec
physical_plan_clickbench_q27                          1.00   1852.4±4.25µs        ? ?/sec    1.00   1855.1±5.01µs        ? ?/sec
physical_plan_clickbench_q28                          1.00      2.2±0.00ms        ? ?/sec    1.00      2.2±0.01ms        ? ?/sec
physical_plan_clickbench_q29                          1.00      2.3±0.00ms        ? ?/sec    1.00      2.3±0.01ms        ? ?/sec
physical_plan_clickbench_q3                           1.00   1626.1±2.99µs        ? ?/sec    1.00   1626.4±5.35µs        ? ?/sec
physical_plan_clickbench_q30                          1.00     15.0±0.02ms        ? ?/sec    1.00     15.1±0.02ms        ? ?/sec
physical_plan_clickbench_q31                          1.00      2.3±0.00ms        ? ?/sec    1.00      2.3±0.00ms        ? ?/sec
physical_plan_clickbench_q32                          1.00      2.3±0.00ms        ? ?/sec    1.01      2.3±0.00ms        ? ?/sec
physical_plan_clickbench_q33                          1.00   1938.8±3.86µs        ? ?/sec    1.00   1938.9±5.90µs        ? ?/sec
physical_plan_clickbench_q34                          1.00   1744.8±4.06µs        ? ?/sec    1.00  1752.3±10.17µs        ? ?/sec
physical_plan_clickbench_q35                          1.00   1785.8±3.46µs        ? ?/sec    1.01  1795.3±17.88µs        ? ?/sec
physical_plan_clickbench_q36                          1.00      2.0±0.00ms        ? ?/sec    1.00      2.0±0.01ms        ? ?/sec
physical_plan_clickbench_q37                          1.00      2.5±0.00ms        ? ?/sec    1.00      2.5±0.01ms        ? ?/sec
physical_plan_clickbench_q38                          1.00      2.5±0.00ms        ? ?/sec    1.00      2.5±0.01ms        ? ?/sec
physical_plan_clickbench_q39                          1.00      2.5±0.00ms        ? ?/sec    1.00      2.5±0.00ms        ? ?/sec
physical_plan_clickbench_q4                           1.00   1466.6±3.87µs        ? ?/sec    1.00   1468.6±4.84µs        ? ?/sec
physical_plan_clickbench_q40                          1.00      3.2±0.01ms        ? ?/sec    1.00      3.2±0.00ms        ? ?/sec
physical_plan_clickbench_q41                          1.00      2.7±0.00ms        ? ?/sec    1.00      2.7±0.01ms        ? ?/sec
physical_plan_clickbench_q42                          1.00      2.8±0.01ms        ? ?/sec    1.00      2.8±0.01ms        ? ?/sec
physical_plan_clickbench_q43                          1.00      3.0±0.01ms        ? ?/sec    1.00      3.0±0.01ms        ? ?/sec
physical_plan_clickbench_q44                          1.00   1540.9±4.01µs        ? ?/sec    1.00   1535.0±5.24µs        ? ?/sec
physical_plan_clickbench_q45                          1.00   1547.9±3.53µs        ? ?/sec    1.00   1540.7±4.94µs        ? ?/sec
physical_plan_clickbench_q46                          1.01   1789.8±4.06µs        ? ?/sec    1.00   1779.1±5.50µs        ? ?/sec
physical_plan_clickbench_q47                          1.00      2.4±0.00ms        ? ?/sec    1.00      2.4±0.01ms        ? ?/sec
physical_plan_clickbench_q48                          1.00      2.5±0.00ms        ? ?/sec    1.00      2.5±0.01ms        ? ?/sec
physical_plan_clickbench_q49                          1.00      2.6±0.00ms        ? ?/sec    1.00      2.6±0.01ms        ? ?/sec
physical_plan_clickbench_q5                           1.00   1582.6±5.10µs        ? ?/sec    1.00   1579.8±4.77µs        ? ?/sec
physical_plan_clickbench_q50                          1.00      2.6±0.00ms        ? ?/sec    1.00      2.6±0.00ms        ? ?/sec
physical_plan_clickbench_q51                          1.01  1903.8±11.68µs        ? ?/sec    1.00   1891.7±6.37µs        ? ?/sec
physical_plan_clickbench_q52                          1.00      2.4±0.01ms        ? ?/sec    1.00      2.4±0.01ms        ? ?/sec
physical_plan_clickbench_q53                          1.01   1750.5±3.91µs        ? ?/sec    1.00   1741.1±9.76µs        ? ?/sec
physical_plan_clickbench_q54                          1.00   1752.9±4.50µs        ? ?/sec    1.00  1748.8±22.14µs        ? ?/sec
physical_plan_clickbench_q55                          1.00   1691.5±4.93µs        ? ?/sec    1.00  1690.0±10.90µs        ? ?/sec
physical_plan_clickbench_q56                          1.00   1692.4±4.66µs        ? ?/sec    1.00   1689.4±7.55µs        ? ?/sec
physical_plan_clickbench_q57                          1.00   1792.2±5.32µs        ? ?/sec    1.01  1806.6±14.98µs        ? ?/sec
physical_plan_clickbench_q58                          1.00      2.2±0.00ms        ? ?/sec    1.02      2.3±0.02ms        ? ?/sec
physical_plan_clickbench_q59                          1.00   1694.4±3.55µs        ? ?/sec    1.01  1709.0±10.17µs        ? ?/sec
physical_plan_clickbench_q6                           1.00   1589.7±4.36µs        ? ?/sec    1.00   1583.6±5.43µs        ? ?/sec
physical_plan_clickbench_q60                          1.00   1695.2±4.50µs        ? ?/sec    1.00   1694.3±6.39µs        ? ?/sec
physical_plan_clickbench_q7                           1.00   1413.3±4.17µs        ? ?/sec    1.10   1560.2±5.57µs        ? ?/sec
physical_plan_clickbench_q8                           1.00   1908.4±3.87µs        ? ?/sec    1.00   1910.6±5.02µs        ? ?/sec
physical_plan_clickbench_q9                           1.00   1856.9±3.91µs        ? ?/sec    1.00   1858.2±6.51µs        ? ?/sec
physical_plan_struct_join_agg_sort                    1.00   1250.5±2.06µs        ? ?/sec    1.02   1273.4±2.72µs        ? ?/sec
physical_plan_tpcds_all                               1.00    655.0±3.26ms        ? ?/sec    1.00    653.1±0.44ms        ? ?/sec
physical_plan_tpch_all                                1.00     42.0±0.03ms        ? ?/sec    1.00     42.1±0.04ms        ? ?/sec
physical_plan_tpch_q1                                 1.00   1417.9±2.93µs        ? ?/sec    1.01   1432.8±3.46µs        ? ?/sec
physical_plan_tpch_q10                                1.00      2.3±0.00ms        ? ?/sec    1.00      2.3±0.00ms        ? ?/sec
physical_plan_tpch_q11                                1.00      2.2±0.00ms        ? ?/sec    1.00      2.2±0.00ms        ? ?/sec
physical_plan_tpch_q12                                1.00   1215.5±1.90µs        ? ?/sec    1.01   1223.4±2.12µs        ? ?/sec
physical_plan_tpch_q13                                1.00   1001.1±1.97µs        ? ?/sec    1.01   1013.5±2.41µs        ? ?/sec
physical_plan_tpch_q14                                1.01   1319.9±2.99µs        ? ?/sec    1.00   1313.3±5.54µs        ? ?/sec
physical_plan_tpch_q16                                1.00   1592.7±1.90µs        ? ?/sec    1.01   1602.3±3.48µs        ? ?/sec
physical_plan_tpch_q17                                1.00   1595.3±2.37µs        ? ?/sec    1.00   1595.5±4.78µs        ? ?/sec
physical_plan_tpch_q18                                1.00   1888.0±2.58µs        ? ?/sec    1.01   1898.6±2.65µs        ? ?/sec
physical_plan_tpch_q19                                1.00   1798.9±3.14µs        ? ?/sec    1.00   1799.2±2.50µs        ? ?/sec
physical_plan_tpch_q2                                 1.00      3.5±0.00ms        ? ?/sec    1.01      3.5±0.00ms        ? ?/sec
physical_plan_tpch_q20                                1.00      2.0±0.00ms        ? ?/sec    1.01      2.1±0.00ms        ? ?/sec
physical_plan_tpch_q21                                1.00      2.7±0.00ms        ? ?/sec    1.00      2.7±0.01ms        ? ?/sec
physical_plan_tpch_q22                                1.00   1473.9±2.35µs        ? ?/sec    1.01   1482.8±5.17µs        ? ?/sec
physical_plan_tpch_q3                                 1.00   1753.2±3.04µs        ? ?/sec    1.01   1769.4±3.26µs        ? ?/sec
physical_plan_tpch_q4                                 1.00   1110.5±1.67µs        ? ?/sec    1.01   1118.2±2.63µs        ? ?/sec
physical_plan_tpch_q5                                 1.00      2.5±0.00ms        ? ?/sec    1.00      2.5±0.00ms        ? ?/sec
physical_plan_tpch_q6                                 1.01    584.3±1.23µs        ? ?/sec    1.00    579.6±1.67µs        ? ?/sec
physical_plan_tpch_q7                                 1.00      2.7±0.00ms        ? ?/sec    1.00      2.7±0.00ms        ? ?/sec
physical_plan_tpch_q8                                 1.00      3.6±0.00ms        ? ?/sec    1.01      3.7±0.01ms        ? ?/sec
physical_plan_tpch_q9                                 1.00      2.6±0.00ms        ? ?/sec    1.00      2.6±0.00ms        ? ?/sec
physical_select_aggregates_from_200                   1.01     12.4±0.02ms        ? ?/sec    1.00     12.2±0.02ms        ? ?/sec
physical_select_all_from_1000                         1.00     19.7±0.05ms        ? ?/sec    1.00     19.8±0.06ms        ? ?/sec
physical_select_one_from_700                          1.00    745.5±2.24µs        ? ?/sec    1.01    755.0±3.18µs        ? ?/sec
physical_sorted_union_order_by_10_int64               1.00      3.8±0.00ms        ? ?/sec    1.01      3.8±0.00ms        ? ?/sec
physical_sorted_union_order_by_10_uint64              1.01      7.5±0.01ms        ? ?/sec    1.00      7.5±0.03ms        ? ?/sec
physical_sorted_union_order_by_50_int64               1.00     83.0±0.26ms        ? ?/sec    1.01     83.5±0.19ms        ? ?/sec
physical_sorted_union_order_by_50_uint64              1.00    294.8±0.45ms        ? ?/sec    1.01    298.6±3.95ms        ? ?/sec
physical_theta_join_consider_sort                     1.00   1057.2±1.49µs        ? ?/sec    1.01   1065.4±2.17µs        ? ?/sec
physical_unnest_to_join                               1.00    618.7±1.92µs        ? ?/sec    1.02    628.2±2.07µs        ? ?/sec
physical_window_function_partition_by_12_on_values    1.00    656.8±1.41µs        ? ?/sec    1.01    660.5±1.23µs        ? ?/sec
physical_window_function_partition_by_30_on_values    1.00   1302.8±1.63µs        ? ?/sec    1.01   1312.7±2.23µs        ? ?/sec
physical_window_function_partition_by_4_on_values     1.03    411.6±1.29µs        ? ?/sec    1.00    401.3±0.71µs        ? ?/sec
physical_window_function_partition_by_7_on_values     1.00    495.7±0.98µs        ? ?/sec    1.00    496.3±1.54µs        ? ?/sec
physical_window_function_partition_by_8_on_values     1.00    532.7±0.86µs        ? ?/sec    1.00    531.3±1.33µs        ? ?/sec
with_param_values_many_columns                        1.00    433.2±2.64µs        ? ?/sec    1.00    434.2±3.26µs        ? ?/sec

Resource Usage

sql_planner — base (merge-base)

Metric Value
Wall time 2060.4s
Peak memory 137.8 MiB
Avg memory 87.5 MiB
CPU user 1984.7s
CPU sys 1.3s
Peak spill 0 B

sql_planner — branch

Metric Value
Wall time 2045.5s
Peak memory 134.1 MiB
Avg memory 86.8 MiB
CPU user 1982.9s
CPU sys 1.4s
Peak spill 0 B

File an issue against this benchmark runner

@alamb

alamb commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

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.

So one way to bring back performance then is to update JoinSelection so it doesn't have to call enforce_distribution_requirements (it can fix the distribution internally if it changes the plan tree).

Sorry I find it really hard to read large / wall of comments, so you may have already said this farther down

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.

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 JoinSelection as a PhysicalAnalyzerRule because it produces invalid plans we could do that initially (and keep the same effective ordering of passes as today, but split between Analyzer/Optimizer)

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.
@zhuqi-lucas

zhuqi-lucas commented Oct 2, 2026 •

Copy link
Copy Markdown
Contributor Author

Thanks @alamb, and sorry for the long comments, I will keep them short from now on.

For example, if we had to initially treat JoinSelection as a PhysicalAnalyzerRule because it produces invalid plans we could do that initially (and keep the same effective ordering of passes as today, but split between Analyzer/Optimizer)

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)

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 EnsureRequirements, in the same order. No rule logic changes and no plan changes at all.

And the benchmark is good, no regression:
#25688 (comment)

The follow-ups are listed in the PR description: moving aggregate_statistics / LimitedDistinctAggregation / FilterPushdown into the optimizer, then WindowTopN, then join_selection (the design discussion you started), and the debug-only boundary check with check_invariants.

@zhuqi-lucas
zhuqi-lucas requested a review from alamb October 2, 2026 06:07
@zhuqi-lucas
zhuqi-lucas force-pushed the physical-analyzer-rule branch from eb879ec to 86fde11 Compare October 2, 2026 06:20
@zhuqi-lucas

Copy link
Copy Markdown
Contributor Author

run benchmark sql_planner

@github-actions github-actions Bot removed sqllogictest SQL Logic Tests (.slt) physical-plan Changes to the physical-plan crate auto detected api change Auto detected API change labels Oct 2, 2026
@adriangbot

Copy link
Copy Markdown

🤖 Benchmark running (GKE) | trigger
Instance: c4a-highcpu-48 (12 vCPU / 65 GiB) | Linux bench-c5946668057-3020-b7dnt 6.12.94+ #1 SMP Wed Aug 19 07:47:20 UTC 2026 aarch64 GNU/Linux

CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  48
On-line CPU(s) list:                     0-47
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     48
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               3 MiB (48 instances)
L1i cache:                               3 MiB (48 instances)
L2 cache:                                96 MiB (48 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-47
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected

Comparing physical-analyzer-rule (79c63da) to 4d167a1 (merge-base) diff

Run configuration
run benchmark sql_planner

Results will be posted here when complete


File an issue against this benchmark runner

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark completed (GKE) | trigger

Instance: c4a-highcpu-48 (12 vCPU / 65 GiB)

Comparing physical-analyzer-rule (79c63da) to 4d167a1 (merge-base) diff

Run configuration
run benchmark sql_planner
CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  48
On-line CPU(s) list:                     0-47
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     48
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               3 MiB (48 instances)
L1i cache:                               3 MiB (48 instances)
L2 cache:                                96 MiB (48 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-47
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected
Details

group                                                 HEAD                                   physical-analyzer-rule
-----                                                 ----                                   ----------------------
logical_aggregate_with_join                           1.00    401.8±1.66µs        ? ?/sec    1.00    400.4±1.66µs        ? ?/sec
logical_correlated_subquery_exists                    1.01    246.5±0.75µs        ? ?/sec    1.00    244.6±0.79µs        ? ?/sec
logical_correlated_subquery_in                        1.01    248.3±1.47µs        ? ?/sec    1.00    246.5±1.62µs        ? ?/sec
logical_distinct_many_columns                         1.00    333.1±0.63µs        ? ?/sec    1.00    331.7±1.29µs        ? ?/sec
logical_join_4_with_agg_and_filter                    1.01    201.2±1.45µs        ? ?/sec    1.00    199.4±0.67µs        ? ?/sec
logical_join_8_with_agg_sort_limit                    1.02    365.4±5.87µs        ? ?/sec    1.00    357.2±1.61µs        ? ?/sec
logical_join_chain_16                                 1.03   653.4±27.64µs        ? ?/sec    1.00    633.0±1.22µs        ? ?/sec
logical_join_chain_4                                  1.03     80.3±0.83µs        ? ?/sec    1.00     78.3±0.23µs        ? ?/sec
logical_join_chain_8                                  1.03    213.6±5.53µs        ? ?/sec    1.00    207.3±0.58µs        ? ?/sec
logical_multiple_subqueries                           1.00    480.1±1.94µs        ? ?/sec    1.02    489.3±1.84µs        ? ?/sec
logical_nested_cte_4_levels                           1.02    222.3±1.11µs        ? ?/sec    1.00    218.6±1.45µs        ? ?/sec
logical_plan_struct_join_agg_sort                     1.01    126.7±0.65µs        ? ?/sec    1.00    126.0±0.70µs        ? ?/sec
logical_plan_tpcds_all                                1.01     83.1±0.20ms        ? ?/sec    1.00     82.0±0.18ms        ? ?/sec
logical_plan_tpch_all                                 1.01      5.6±0.03ms        ? ?/sec    1.00      5.6±0.01ms        ? ?/sec
logical_scalar_subquery                               1.00    266.1±0.78µs        ? ?/sec    1.00    265.0±0.86µs        ? ?/sec
logical_select_all_from_1000                          1.00      9.1±0.03ms        ? ?/sec    1.01      9.2±0.04ms        ? ?/sec
logical_select_one_from_700                           1.00    290.5±1.43µs        ? ?/sec    1.02    295.3±2.00µs        ? ?/sec
logical_trivial_join_high_numbered_columns            1.00    246.9±1.45µs        ? ?/sec    1.00    247.7±1.75µs        ? ?/sec
logical_trivial_join_low_numbered_columns             1.00    234.0±1.34µs        ? ?/sec    1.00    234.4±1.57µs        ? ?/sec
logical_union_4_branches                              1.00    391.5±1.78µs        ? ?/sec    1.00    392.2±7.14µs        ? ?/sec
logical_union_8_branches                              1.00    792.7±1.91µs        ? ?/sec    1.00   789.7±13.04µs        ? ?/sec
logical_wide_aggregate_1000_exprs                     1.01     44.3±0.12ms        ? ?/sec    1.00     44.0±0.08ms        ? ?/sec
logical_wide_aggregate_100_exprs                      1.00   1640.7±2.94µs        ? ?/sec    1.00   1638.6±2.89µs        ? ?/sec
logical_wide_case_50_exprs                            1.01   1816.5±4.75µs        ? ?/sec    1.00   1807.3±4.74µs        ? ?/sec
logical_wide_filter_200_predicates                    1.01   1336.0±5.98µs        ? ?/sec    1.00   1325.3±6.68µs        ? ?/sec
logical_wide_filter_50_predicates                     1.02    368.1±2.03µs        ? ?/sec    1.00    360.6±1.48µs        ? ?/sec
optimizer_correlated_exists                           1.00    224.4±0.66µs        ? ?/sec    1.00    223.9±0.49µs        ? ?/sec
optimizer_join_4_with_agg_filter                      1.02    422.4±0.77µs        ? ?/sec    1.00    413.7±1.40µs        ? ?/sec
optimizer_join_chain_4                                1.02    163.9±0.27µs        ? ?/sec    1.00    161.4±0.53µs        ? ?/sec
optimizer_join_chain_8                                1.00    550.3±0.79µs        ? ?/sec    1.00    549.3±0.87µs        ? ?/sec
optimizer_select_all_from_1000                        1.00      7.1±0.04ms        ? ?/sec    1.00      7.1±0.02ms        ? ?/sec
optimizer_select_one_from_700                         1.00    233.5±0.34µs        ? ?/sec    1.01    235.3±0.52µs        ? ?/sec
optimizer_tpcds_all                                   1.03    291.2±2.55ms        ? ?/sec    1.00    284.0±0.42ms        ? ?/sec
optimizer_tpch_all                                    1.03     16.0±0.06ms        ? ?/sec    1.00     15.5±0.02ms        ? ?/sec
optimizer_wide_aggregate_100                          1.01   1800.2±3.04µs        ? ?/sec    1.00   1781.1±3.73µs        ? ?/sec
optimizer_wide_filter_200                             1.01      4.0±0.01ms        ? ?/sec    1.00      3.9±0.01ms        ? ?/sec
physical_intersection                                 1.01    574.6±2.43µs        ? ?/sec    1.00    567.8±2.42µs        ? ?/sec
physical_join_consider_sort                           1.00   1039.6±7.16µs        ? ?/sec    1.00   1035.6±5.31µs        ? ?/sec
physical_join_distinct                                1.01    227.7±1.50µs        ? ?/sec    1.00    226.4±1.52µs        ? ?/sec
physical_many_self_joins                              1.00      7.6±0.01ms        ? ?/sec    1.01      7.7±0.03ms        ? ?/sec
physical_plan_clickbench_all                          1.01    139.4±0.37ms        ? ?/sec    1.00    138.1±0.26ms        ? ?/sec
physical_plan_clickbench_q1                           1.00   1387.7±7.54µs        ? ?/sec    1.00   1386.1±8.65µs        ? ?/sec
physical_plan_clickbench_q10                          1.00      2.0±0.01ms        ? ?/sec    1.00      2.0±0.01ms        ? ?/sec
physical_plan_clickbench_q11                          1.01      2.2±0.01ms        ? ?/sec    1.00      2.2±0.01ms        ? ?/sec
physical_plan_clickbench_q12                          1.00      2.2±0.01ms        ? ?/sec    1.00      2.3±0.01ms        ? ?/sec
physical_plan_clickbench_q13                          1.00      2.0±0.01ms        ? ?/sec    1.00      2.0±0.01ms        ? ?/sec
physical_plan_clickbench_q14                          1.00      2.2±0.01ms        ? ?/sec    1.00      2.2±0.01ms        ? ?/sec
physical_plan_clickbench_q15                          1.00      2.1±0.01ms        ? ?/sec    1.00      2.1±0.01ms        ? ?/sec
physical_plan_clickbench_q16                          1.00   1796.3±6.09µs        ? ?/sec    1.00   1794.8±5.47µs        ? ?/sec
physical_plan_clickbench_q17                          1.00   1845.0±5.95µs        ? ?/sec    1.00   1841.7±5.70µs        ? ?/sec
physical_plan_clickbench_q18                          1.00   1698.2±6.22µs        ? ?/sec    1.00   1690.8±4.28µs        ? ?/sec
physical_plan_clickbench_q19                          1.00      2.0±0.01ms        ? ?/sec    1.00      2.0±0.01ms        ? ?/sec
physical_plan_clickbench_q2                           1.02   1789.5±7.01µs        ? ?/sec    1.00   1762.5±4.39µs        ? ?/sec
physical_plan_clickbench_q20                          1.00   1569.2±6.26µs        ? ?/sec    1.00   1564.5±4.68µs        ? ?/sec
physical_plan_clickbench_q21                          1.00   1783.9±7.80µs        ? ?/sec    1.01  1794.4±13.25µs        ? ?/sec
physical_plan_clickbench_q22                          1.00      2.1±0.01ms        ? ?/sec    1.00      2.1±0.01ms        ? ?/sec
physical_plan_clickbench_q23                          1.00      2.3±0.01ms        ? ?/sec    1.00      2.3±0.01ms        ? ?/sec
physical_plan_clickbench_q24                          1.00      5.2±0.01ms        ? ?/sec    1.00      5.2±0.01ms        ? ?/sec
physical_plan_clickbench_q25                          1.00   1880.9±6.59µs        ? ?/sec    1.00   1882.1±5.50µs        ? ?/sec
physical_plan_clickbench_q26                          1.00   1744.0±6.39µs        ? ?/sec    1.00   1744.3±5.44µs        ? ?/sec
physical_plan_clickbench_q27                          1.00   1912.0±6.29µs        ? ?/sec    1.00   1914.5±5.39µs        ? ?/sec
physical_plan_clickbench_q28                          1.00      2.3±0.01ms        ? ?/sec    1.00      2.3±0.00ms        ? ?/sec
physical_plan_clickbench_q29                          1.00      2.4±0.01ms        ? ?/sec    1.00      2.4±0.01ms        ? ?/sec
physical_plan_clickbench_q3                           1.01   1675.3±6.90µs        ? ?/sec    1.00   1651.5±7.65µs        ? ?/sec
physical_plan_clickbench_q30                          1.00     12.8±0.02ms        ? ?/sec    1.00     12.8±0.02ms        ? ?/sec
physical_plan_clickbench_q31                          1.00      2.4±0.01ms        ? ?/sec    1.00      2.4±0.01ms        ? ?/sec
physical_plan_clickbench_q32                          1.00      2.4±0.01ms        ? ?/sec    1.00      2.4±0.01ms        ? ?/sec
physical_plan_clickbench_q33                          1.00      2.0±0.00ms        ? ?/sec    1.00   1998.5±3.97µs        ? ?/sec
physical_plan_clickbench_q34                          1.00   1803.2±6.75µs        ? ?/sec    1.00   1809.2±4.56µs        ? ?/sec
physical_plan_clickbench_q35                          1.00   1847.0±6.26µs        ? ?/sec    1.00   1849.6±5.89µs        ? ?/sec
physical_plan_clickbench_q36                          1.00      2.1±0.01ms        ? ?/sec    1.00      2.1±0.01ms        ? ?/sec
physical_plan_clickbench_q37                          1.01      2.5±0.01ms        ? ?/sec    1.00      2.5±0.01ms        ? ?/sec
physical_plan_clickbench_q38                          1.01      2.5±0.01ms        ? ?/sec    1.00      2.5±0.00ms        ? ?/sec
physical_plan_clickbench_q39                          1.02      2.6±0.01ms        ? ?/sec    1.00      2.6±0.01ms        ? ?/sec
physical_plan_clickbench_q4                           1.02   1521.7±6.61µs        ? ?/sec    1.00   1492.0±4.35µs        ? ?/sec
physical_plan_clickbench_q40                          1.00      3.2±0.01ms        ? ?/sec    1.00      3.2±0.01ms        ? ?/sec
physical_plan_clickbench_q41                          1.00      2.8±0.01ms        ? ?/sec    1.00      2.8±0.01ms        ? ?/sec
physical_plan_clickbench_q42                          1.01      2.9±0.01ms        ? ?/sec    1.00      2.9±0.01ms        ? ?/sec
physical_plan_clickbench_q43                          1.01      3.1±0.01ms        ? ?/sec    1.00      3.1±0.01ms        ? ?/sec
physical_plan_clickbench_q44                          1.01   1593.8±5.82µs        ? ?/sec    1.00   1577.8±4.56µs        ? ?/sec
physical_plan_clickbench_q45                          1.01   1600.7±6.42µs        ? ?/sec    1.00   1582.4±5.26µs        ? ?/sec
physical_plan_clickbench_q46                          1.01   1846.7±6.94µs        ? ?/sec    1.00   1831.6±6.23µs        ? ?/sec
physical_plan_clickbench_q47                          1.01      2.4±0.01ms        ? ?/sec    1.00      2.4±0.00ms        ? ?/sec
physical_plan_clickbench_q48                          1.01      2.6±0.01ms        ? ?/sec    1.00      2.6±0.00ms        ? ?/sec
physical_plan_clickbench_q49                          1.01      2.6±0.01ms        ? ?/sec    1.00      2.6±0.01ms        ? ?/sec
physical_plan_clickbench_q5                           1.01   1636.6±6.11µs        ? ?/sec    1.00   1613.1±3.82µs        ? ?/sec
physical_plan_clickbench_q50                          1.00      2.7±0.01ms        ? ?/sec    1.01      2.7±0.02ms        ? ?/sec
physical_plan_clickbench_q51                          1.01   1955.0±7.59µs        ? ?/sec    1.00   1933.4±5.38µs        ? ?/sec
physical_plan_clickbench_q52                          1.01      2.5±0.01ms        ? ?/sec    1.00      2.4±0.01ms        ? ?/sec
physical_plan_clickbench_q53                          1.01   1806.0±6.83µs        ? ?/sec    1.00   1787.4±4.09µs        ? ?/sec
physical_plan_clickbench_q54                          1.01   1808.0±7.60µs        ? ?/sec    1.00   1787.8±5.01µs        ? ?/sec
physical_plan_clickbench_q55                          1.01   1747.4±6.79µs        ? ?/sec    1.00   1724.2±4.29µs        ? ?/sec
physical_plan_clickbench_q56                          1.02   1750.4±7.16µs        ? ?/sec    1.00   1723.9±4.13µs        ? ?/sec
physical_plan_clickbench_q57                          1.02   1856.5±6.63µs        ? ?/sec    1.00   1827.8±4.80µs        ? ?/sec
physical_plan_clickbench_q58                          1.01      2.3±0.01ms        ? ?/sec    1.00      2.3±0.01ms        ? ?/sec
physical_plan_clickbench_q59                          1.01   1752.8±6.28µs        ? ?/sec    1.00   1729.1±4.58µs        ? ?/sec
physical_plan_clickbench_q6                           1.02   1641.3±6.32µs        ? ?/sec    1.00   1615.5±4.19µs        ? ?/sec
physical_plan_clickbench_q60                          1.01   1753.3±6.81µs        ? ?/sec    1.00   1727.9±3.52µs        ? ?/sec
physical_plan_clickbench_q7                           1.02   1458.4±6.38µs        ? ?/sec    1.00   1431.6±3.62µs        ? ?/sec
physical_plan_clickbench_q8                           1.02   1975.7±8.40µs        ? ?/sec    1.00   1942.6±4.92µs        ? ?/sec
physical_plan_clickbench_q9                           1.02   1919.2±6.99µs        ? ?/sec    1.00   1890.1±4.90µs        ? ?/sec
physical_plan_struct_join_agg_sort                    1.00   1269.7±2.66µs        ? ?/sec    1.00   1268.3±3.39µs        ? ?/sec
physical_plan_tpcds_all                               1.00    657.2±0.51ms        ? ?/sec    1.00    654.7±1.19ms        ? ?/sec
physical_plan_tpch_all                                1.00     42.7±0.04ms        ? ?/sec    1.00     42.7±0.31ms        ? ?/sec
physical_plan_tpch_q1                                 1.01   1420.6±2.43µs        ? ?/sec    1.00   1410.5±4.10µs        ? ?/sec
physical_plan_tpch_q10                                1.00      2.4±0.00ms        ? ?/sec    1.00      2.3±0.00ms        ? ?/sec
physical_plan_tpch_q11                                1.01      2.2±0.00ms        ? ?/sec    1.00      2.2±0.00ms        ? ?/sec
physical_plan_tpch_q12                                1.01   1249.4±4.02µs        ? ?/sec    1.00   1235.0±1.78µs        ? ?/sec
physical_plan_tpch_q13                                1.03   1052.6±1.59µs        ? ?/sec    1.00   1026.0±2.28µs        ? ?/sec
physical_plan_tpch_q14                                1.01   1345.1±4.55µs        ? ?/sec    1.00   1333.8±2.51µs        ? ?/sec
physical_plan_tpch_q16                                1.01   1630.3±3.82µs        ? ?/sec    1.00   1613.2±2.07µs        ? ?/sec
physical_plan_tpch_q17                                1.01   1628.2±4.08µs        ? ?/sec    1.00   1614.3±3.40µs        ? ?/sec
physical_plan_tpch_q18                                1.01   1912.9±3.85µs        ? ?/sec    1.00   1899.5±2.69µs        ? ?/sec
physical_plan_tpch_q19                                1.01   1832.8±4.31µs        ? ?/sec    1.00   1813.8±3.52µs        ? ?/sec
physical_plan_tpch_q2                                 1.00      3.5±0.00ms        ? ?/sec    1.00      3.5±0.01ms        ? ?/sec
physical_plan_tpch_q20                                1.01      2.1±0.00ms        ? ?/sec    1.00      2.1±0.00ms        ? ?/sec
physical_plan_tpch_q21                                1.01      2.7±0.00ms        ? ?/sec    1.00      2.7±0.00ms        ? ?/sec
physical_plan_tpch_q22                                1.01   1518.4±3.91µs        ? ?/sec    1.00   1499.8±2.44µs        ? ?/sec
physical_plan_tpch_q3                                 1.01   1787.8±3.75µs        ? ?/sec    1.00   1764.2±3.11µs        ? ?/sec
physical_plan_tpch_q4                                 1.01   1145.8±1.81µs        ? ?/sec    1.00   1130.8±2.74µs        ? ?/sec
physical_plan_tpch_q5                                 1.00      2.5±0.00ms        ? ?/sec    1.00      2.5±0.00ms        ? ?/sec
physical_plan_tpch_q6                                 1.02    596.6±2.02µs        ? ?/sec    1.00    585.7±1.43µs        ? ?/sec
physical_plan_tpch_q7                                 1.00      2.7±0.00ms        ? ?/sec    1.00      2.7±0.00ms        ? ?/sec
physical_plan_tpch_q8                                 1.00      3.7±0.01ms        ? ?/sec    1.00      3.6±0.00ms        ? ?/sec
physical_plan_tpch_q9                                 1.00      2.6±0.00ms        ? ?/sec    1.00      2.6±0.00ms        ? ?/sec
physical_select_aggregates_from_200                   1.01      7.9±0.01ms        ? ?/sec    1.00      7.8±0.02ms        ? ?/sec
physical_select_all_from_1000                         1.00     20.6±0.07ms        ? ?/sec    1.00     20.6±0.08ms        ? ?/sec
physical_select_one_from_700                          1.00    736.8±3.15µs        ? ?/sec    1.00    738.5±2.45µs        ? ?/sec
physical_sorted_union_order_by_10_int64               1.00      3.8±0.00ms        ? ?/sec    1.03      3.9±0.04ms        ? ?/sec
physical_sorted_union_order_by_10_uint64              1.00      7.5±0.01ms        ? ?/sec    1.00      7.5±0.01ms        ? ?/sec
physical_sorted_union_order_by_50_int64               1.00     83.5±0.22ms        ? ?/sec    1.01     84.6±0.27ms        ? ?/sec
physical_sorted_union_order_by_50_uint64              1.00    297.4±0.53ms        ? ?/sec    1.02    304.3±2.35ms        ? ?/sec
physical_theta_join_consider_sort                     1.00   1065.5±7.30µs        ? ?/sec    1.00   1062.3±5.44µs        ? ?/sec
physical_unnest_to_join                               1.01    621.1±1.84µs        ? ?/sec    1.00    615.9±2.05µs        ? ?/sec
physical_window_function_partition_by_12_on_values    1.01    674.5±6.15µs        ? ?/sec    1.00    667.3±1.69µs        ? ?/sec
physical_window_function_partition_by_30_on_values    1.00   1321.5±3.63µs        ? ?/sec    1.00   1315.7±3.91µs        ? ?/sec
physical_window_function_partition_by_4_on_values     1.01    411.9±1.07µs        ? ?/sec    1.00    406.5±1.73µs        ? ?/sec
physical_window_function_partition_by_7_on_values     1.01    506.0±1.18µs        ? ?/sec    1.00    501.5±1.23µs        ? ?/sec
physical_window_function_partition_by_8_on_values     1.02    546.5±5.40µs        ? ?/sec    1.00    538.2±1.11µs        ? ?/sec
with_param_values_many_columns                        1.02    444.2±2.95µs        ? ?/sec    1.00    437.5±1.94µs        ? ?/sec

Resource Usage

sql_planner — base (merge-base)

Metric Value
Wall time 2060.5s
Peak memory 136.7 MiB
Avg memory 88.9 MiB
CPU user 2014.1s
CPU sys 1.4s
Peak spill 0 B

sql_planner — branch

Metric Value
Wall time 2055.5s
Peak memory 136.9 MiB
Avg memory 89.1 MiB
CPU user 2005.2s
CPU sys 1.5s
Peak spill 0 B

File an issue against this benchmark runner

@2010YOUY01

2010YOUY01 commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

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 implementation

The latest change keeps the existing physical optimizer rule order, but introduces one hard boundary through the PhysicalAnalyzer / PhysicalOptimizer split.

Conceptually:

  • Before the boundary: there are no enforced repartition/sort operators yet, so rules are relatively free to transform the plan.
  • After the boundary: repartition/sort operators may already exist, so every later rule must recognize and preserve the distribution and ordering requirements that have already been satisfied.

My main suggestion: Extensible Boundaries

It would be very helpful to make this boundary mechanism extensible, instead of encoding one specific boundary through the Analyzer/Optimizer split.

# Current design

PhysicalAnalyzer:
    OutputRequirements(add)
    aggregate_statistics
    join_selection
    LimitedDistinctAggregation
    FilterPushdown
    WindowTopN
    EnsureRequirements

    # Boundary:
    # distribution and ordering requirements are satisfied

PhysicalOptimizer:
    ...remaining existing rules...


# What I mean by extensible boundaries

rule1
rule2
rule3

    # Boundary A:
    # distribution and ordering requirements are now enforced

rule4
rule5

    # Boundary B:
    # another optimizer invariant is now established

rule6
rule7

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:

rule1

    # Before this point, projection has one canonical representation:
    #
    # ProjectionExec
    # -- child

ProjectionPushDown

    # After this point, projection may either stay explicit:
    #
    # ProjectionExec(c1 + c2)
    # -- Join(output = [c1, c2])
    #
    # or be fused into an operator that supports it:
    #
    # Join(
    #     output = [c1, c2],
    #     projection = [c1],
    # )

rule2
rule3

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 boundaries

I 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:

EnsureRequirements

Rule3
ValidateEnsureRequirements

Rule4
ValidateEnsureRequirements

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.

@zhuqi-lucas

Copy link
Copy Markdown
Contributor Author

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 cfg(debug_assertions), run the sanity check (distribution and ordering satisfied) after every optimizer rule, so a violation points at the rule that caused it. Zero cost in release, same as the existing per-rule check_invariants. Happy to add that here if you both prefer, otherwise as a follow-up.

@alamb

alamb commented Oct 6, 2026 •

Copy link
Copy Markdown
Contributor

What I mean by extensible boundaries

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) 🤔

@alamb

alamb commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

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

@2010YOUY01

2010YOUY01 commented Oct 6, 2026 •

Copy link
Copy Markdown
Contributor

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 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.)


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

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 extensions

If we explicitly exclude arbitrary reordering of the default optimizer rule list from the supported use case, then stronger constraints can simplify the optimizer:

  • Simpler implementation: each optimizer rule needs to handle fewer possible input states if its preconditions are known.
  • Safer composition: a rule knows which invariants it must preserve so that later rules continue to work correctly.

Walkthrough

In @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 extension_rule become much clearer:

  • Assumption: input plans already satisfy their distribution/ordering requirements (RepartitionExecs have been inserted and the plan is runnable).
  • Promise: the rule must not invalidate those requirements. It either preserves the existing enforcement operators or re-runs EnsureRequirements internally when necessary.

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 PR

Now @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.

@2010YOUY01

Copy link
Copy Markdown
Contributor

If I understand the assumptions behind the optimizer implementation correctly, adding more explicit constraints should make both understanding and implementation easier.

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.

@alamb

alamb commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

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.

Now @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.

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:

  1. Schema matches
  2. Declared required properties match actual input properties

For example it is not clear to me the property described after LimitedDistinctAggregation in https://github.com/apache/datafusion/pull/26069/changes#r4194462241 is actually invariant won't ever change by all possible future optimizer passes, or even if we should try to enforce that as an invariant after subsequent passes

I think my main concerns are:

  1. Introducing a level of abstraction / API that isn't really used for anything other than the analyzer/optimizer split
  2. Adding unnecessarily strict requirements on other optimizer passes

My current preference is to implement the more general optimizer boundary directly, to avoid introducing multiple mechanisms for the same underlying problem.

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.

If we can come up with some other specific, useful boundaries / invariants this sounds like a good plan to me.

@alamb

alamb commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

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.

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

@zhuqi-lucas

Copy link
Copy Markdown
Contributor Author

Thank you @alamb and @2010YOUY01.

One observation from reading #26069 against the code: its "only shape 1 is possible" precondition holds exactly because LimitedDistinctAggregation runs before EnsureRequirements and CombinePartialFinalAggregate, the only producers of the other two shapes. This PR makes that precondition explicit, since the rule now sits in the analyzer list ahead of enforcement. So I think the two PRs agree with each other.

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 SanityCheckPlan already verifies them. Shape facts like #26069's are a different kind: they describe what earlier rules produced, and nothing in the plan can verify them. For those I think the cheapest honest tool is what #26069 already does, a precondition stated in the rule's docs plus a debug assertion inside the rule, rather than a pipeline-level boundary that every rule has to declare against.

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.

@2010YOUY01

Copy link
Copy Markdown
Contributor

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:

  • How should we enforce those invariants?
  • How far can we go with this approach? Adding more constraints might make future changes harder to make, and even increase maintenance overhead at some point.

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:

  • If the result is strong, I'll share it as a PoC PR.
  • Otherwise, I think the lighter-weight approach @zhuqi-lucas suggested makes sense: document these preconditions locally in the relevant optimizer rules, with debug assertions where practical, rather than introducing a general boundary mechanism.

@alamb

alamb commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

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

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

core Core DataFusion crate documentation Improvements or additions to documentation ffi Changes to the ffi crate optimizer Optimizer rules proto Related to proto crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

6 participants