Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
fdcd277
feat: split physical enforcement into a PhysicalAnalyzerRule phase
zhuqi-lucas Sep 24, 2026
216b9aa
refactor: decompose EnsureRequirements into EnforceDistribution / Enf…
zhuqi-lucas Sep 25, 2026
951b734
feat: run distribution enforcement first, as an analyzer, with self-m…
zhuqi-lucas Sep 26, 2026
9ffd58a
feat: move all requirement enforcement into the physical analyzer
zhuqi-lucas Sep 26, 2026
6ea3200
fix: make OutputRequirementExec execute transparently; pin test targe…
zhuqi-lucas Sep 27, 2026
82f80ae
fix: clear analyzer rules in tests that observe an unoptimized physic…
zhuqi-lucas Sep 27, 2026
51642ae
fix: pin target_partitions in window_topn / join_selection tests
zhuqi-lucas Sep 27, 2026
2e3668a
fix: keep enforcement out of memory_limit scenarios that pin custom r…
zhuqi-lucas Sep 27, 2026
0bf1852
fix: apply prettier formatting to optimizer_rule_reference.md
zhuqi-lucas Sep 27, 2026
f9cb358
docs: address review feedback on the analyzer-phase docs
zhuqi-lucas Sep 27, 2026
64f1efb
Merge remote-tracking branch 'apache/main' into enforce-first-full
zhuqi-lucas Sep 28, 2026
9d47f77
Merge branch 'main' into physical-analyzer-rule
zhuqi-lucas Sep 28, 2026
e0c199a
docs: add a phase diagram for the physical analyzer / optimizer split
zhuqi-lucas Sep 28, 2026
c6f93f6
Merge branch 'main' into physical-analyzer-rule
alamb Oct 1, 2026
20c7daf
Run the rules that decide requirements before enforcing them
zhuqi-lucas Oct 1, 2026
b211ce1
Check the analyzer phase's contract where the phase ends
zhuqi-lucas Oct 1, 2026
67d6f4a
Drop stale comments that still describe enforcement as an optimizer rule
zhuqi-lucas Oct 1, 2026
9a39fe3
Document the physical analyzer phase in the 56.0.0 upgrade guide
zhuqi-lucas Oct 1, 2026
267e70c
Keep EnsureRequirements as the single enforcement rule
zhuqi-lucas Oct 1, 2026
64e8179
Merge branch 'physical-analyzer-rule' of github.com:zhuqi-lucas/arrow…
zhuqi-lucas Oct 1, 2026
bc862ce
Describe the enforcement phases as the two walks the code runs
zhuqi-lucas Oct 2, 2026
86fde11
Reduce the analyzer split to the trait, the plumbing and today's rule…
zhuqi-lucas Oct 2, 2026
79c63da
Merge branch 'main' into physical-analyzer-rule
zhuqi-lucas Oct 2, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
83 changes: 82 additions & 1 deletion datafusion/core/src/execution/session_state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -75,10 +75,13 @@ use datafusion_optimizer::{
};
use datafusion_physical_expr::create_physical_expr;
use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
use datafusion_physical_optimizer::analyzer::PhysicalAnalyzer;
use datafusion_physical_optimizer::optimizer::PhysicalOptimizer;
use datafusion_physical_plan::ExecutionPlan;
use datafusion_physical_plan::operator_statistics::StatisticsRegistry;
use datafusion_session::{PhysicalOptimizerContext, PhysicalOptimizerRule, Session};
use datafusion_session::{
PhysicalAnalyzerRule, PhysicalOptimizerContext, PhysicalOptimizerRule, Session,
};
#[cfg(feature = "sql")]
use datafusion_sql::{
parser::{DFParserBuilder, Statement},
Expand Down Expand Up @@ -168,6 +171,9 @@ struct SessionStateInner {
type_planner: Option<Arc<dyn TypePlanner>>,
/// 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.

❤️

physical_analyzers: PhysicalAnalyzer,
/// Responsible for optimizing a physical execution plan
physical_optimizers: PhysicalOptimizer,
/// Responsible for planning `LogicalPlan`s, and `ExecutionPlan`
Expand Down Expand Up @@ -272,6 +278,7 @@ impl Debug for SessionState {
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

.field("physical_optimizers", &self.inner.physical_optimizers)
.field("table_functions", &self.inner.table_functions)
.field("scalar_functions", &self.inner.scalar_functions)
Expand Down Expand Up @@ -312,6 +319,10 @@ impl Session for SessionState {
SessionState::physical_optimizers(self)
}

fn physical_analyzers(&self) -> &[Arc<dyn PhysicalAnalyzerRule + Send + Sync>] {
SessionState::physical_analyzers(self)
}

fn statistics_registry(&self) -> Option<&StatisticsRegistry> {
SessionState::statistics_registry(self)
}
Expand Down Expand Up @@ -454,6 +465,19 @@ impl SessionState {
self
}

/// Add `physical_analyzer_rule` to the end of the list of
/// [`PhysicalAnalyzerRule`]s run before physical optimization.
pub fn add_physical_analyzer_rule(
&mut self,
physical_analyzer_rule: Arc<dyn PhysicalAnalyzerRule + Send + Sync>,
) -> &Self {
Arc::make_mut(&mut self.inner)
.physical_analyzers
.rules
.push(physical_analyzer_rule);
self
}

// the add_optimizer_rule takes an owned reference
// it should probably be renamed to `with_optimizer_rule` to follow builder style
// and `add_optimizer_rule` that takes &mut self added instead of this
Expand Down Expand Up @@ -946,6 +970,11 @@ impl SessionState {
&self.inner.physical_optimizers.rules
}

/// Return the physical analyzers
pub fn physical_analyzers(&self) -> &[Arc<dyn PhysicalAnalyzerRule + Send + Sync>] {
&self.inner.physical_analyzers.rules
}

/// return the configuration options
pub fn config_options(&self) -> &Arc<ConfigOptions> {
self.inner.config.options()
Expand Down Expand Up @@ -1144,6 +1173,7 @@ pub struct SessionStateBuilder {
#[cfg(feature = "sql")]
type_planner: Option<Arc<dyn TypePlanner>>,
optimizer: Option<Optimizer>,
physical_analyzers: Option<PhysicalAnalyzer>,
physical_optimizers: Option<PhysicalOptimizer>,
query_planner: Option<Arc<dyn QueryPlanner + Send + Sync>>,
catalog_list: Option<Arc<dyn CatalogProviderList>>,
Expand All @@ -1167,6 +1197,7 @@ pub struct SessionStateBuilder {
// fields to support convenience functions
analyzer_rules: Option<Vec<Arc<dyn AnalyzerRule + Send + Sync>>>,
optimizer_rules: Option<Vec<Arc<dyn OptimizerRule + Send + Sync>>>,
physical_analyzer_rules: Option<Vec<Arc<dyn PhysicalAnalyzerRule + Send + Sync>>>,
physical_optimizer_rules: Option<Vec<Arc<dyn PhysicalOptimizerRule + Send + Sync>>>,
}

Expand All @@ -1188,6 +1219,7 @@ impl SessionStateBuilder {
#[cfg(feature = "sql")]
type_planner: None,
optimizer: None,
physical_analyzers: None,
physical_optimizers: None,
query_planner: None,
catalog_list: None,
Expand All @@ -1211,6 +1243,7 @@ impl SessionStateBuilder {
// fields to support convenience functions
analyzer_rules: None,
optimizer_rules: None,
physical_analyzer_rules: None,
physical_optimizer_rules: None,
}
}
Expand Down Expand Up @@ -1250,6 +1283,7 @@ impl SessionStateBuilder {
#[cfg(feature = "sql")]
type_planner: existing.type_planner,
optimizer: Some(existing.optimizer),
physical_analyzers: Some(existing.physical_analyzers),
physical_optimizers: Some(existing.physical_optimizers),
query_planner: Some(existing.query_planner),
catalog_list: Some(existing.catalog_list),
Expand Down Expand Up @@ -1277,6 +1311,7 @@ impl SessionStateBuilder {
// fields to support convenience functions
analyzer_rules: None,
optimizer_rules: None,
physical_analyzer_rules: None,
physical_optimizer_rules: None,
}
}
Expand Down Expand Up @@ -1442,6 +1477,27 @@ impl SessionStateBuilder {
self
}

/// Set the [`PhysicalAnalyzerRule`]s run before physical optimization.
pub fn with_physical_analyzer_rules(
mut self,
physical_analyzers: Vec<Arc<dyn PhysicalAnalyzerRule + Send + Sync>>,
) -> Self {
self.physical_analyzers = Some(PhysicalAnalyzer::with_rules(physical_analyzers));
self
}

/// Add `physical_analyzer_rule` to the end of the list of
/// [`PhysicalAnalyzerRule`]s run before physical optimization.
pub fn with_physical_analyzer_rule(
mut self,
physical_analyzer_rule: Arc<dyn PhysicalAnalyzerRule + Send + Sync>,
) -> Self {
let mut rules = self.physical_analyzer_rules.unwrap_or_default();
rules.push(physical_analyzer_rule);
self.physical_analyzer_rules = Some(rules);
self
}

/// Set the [`QueryPlanner`]
pub fn with_query_planner(
mut self,
Expand Down Expand Up @@ -1689,6 +1745,7 @@ impl SessionStateBuilder {
#[cfg(feature = "sql")]
type_planner,
optimizer,
physical_analyzers,
physical_optimizers,
query_planner,
catalog_list,
Expand All @@ -1711,6 +1768,7 @@ impl SessionStateBuilder {
statistics_registry,
analyzer_rules,
optimizer_rules,
physical_analyzer_rules,
physical_optimizer_rules,
} = self;

Expand All @@ -1726,6 +1784,7 @@ impl SessionStateBuilder {
#[cfg(feature = "sql")]
type_planner,
optimizer: optimizer.unwrap_or_default(),
physical_analyzers: physical_analyzers.unwrap_or_default(),
physical_optimizers: physical_optimizers.unwrap_or_default(),
query_planner: query_planner
.unwrap_or_else(|| Arc::new(DefaultQueryPlanner {})),
Expand Down Expand Up @@ -1871,6 +1930,14 @@ impl SessionStateBuilder {
}
}

if let Some(physical_analyzer_rules) = physical_analyzer_rules {
let physical_analyzers =
&mut Arc::make_mut(&mut state.inner).physical_analyzers;
for physical_analyzer_rule in physical_analyzer_rules {
physical_analyzers.rules.push(physical_analyzer_rule);
}
}

if let Some(physical_optimizer_rules) = physical_optimizer_rules {
let physical_optimizers =
&mut Arc::make_mut(&mut state.inner).physical_optimizers;
Expand Down Expand Up @@ -1919,6 +1986,11 @@ impl SessionStateBuilder {
&mut self.physical_optimizers
}

/// Returns the current physical_analyzers value
pub fn physical_analyzers(&mut self) -> &mut Option<PhysicalAnalyzer> {
&mut self.physical_analyzers
}

/// Returns the current query_planner value
pub fn query_planner(&mut self) -> &mut Option<Arc<dyn QueryPlanner + Send + Sync>> {
&mut self.query_planner
Expand Down Expand Up @@ -2030,6 +2102,13 @@ impl SessionStateBuilder {
) -> &mut Option<Vec<Arc<dyn PhysicalOptimizerRule + Send + Sync>>> {
&mut self.physical_optimizer_rules
}

/// Returns the current physical_analyzer_rules value
pub fn physical_analyzer_rules(
&mut self,
) -> &mut Option<Vec<Arc<dyn PhysicalAnalyzerRule + Send + Sync>>> {
&mut self.physical_analyzer_rules
}
}

impl Debug for SessionStateBuilder {
Expand Down Expand Up @@ -2058,6 +2137,8 @@ impl Debug for SessionStateBuilder {
.field("analyzer", &self.analyzer)
.field("optimizer_rules", &self.optimizer_rules)
.field("optimizer", &self.optimizer)
.field("physical_analyzer_rules", &self.physical_analyzer_rules)
.field("physical_analyzers", &self.physical_analyzers)
.field("physical_optimizer_rules", &self.physical_optimizer_rules)
.field("physical_optimizers", &self.physical_optimizers)
.field("table_functions", &self.table_functions)
Expand Down
Loading
Loading