Repository navigation
feat: run physical requirement enforcement first, as a PhysicalAnalyzer phase #25688
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
zhuqi-lucas
wants to merge
23
commits into
apache:main
Choose a base branch
from
zhuqi-lucas:physical-analyzer-rule
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
+917
−94
Open
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 216b9aa
refactor: decompose EnsureRequirements into EnforceDistribution / Enf…
zhuqi-lucas 951b734
feat: run distribution enforcement first, as an analyzer, with self-m…
zhuqi-lucas 9ffd58a
feat: move all requirement enforcement into the physical analyzer
zhuqi-lucas 6ea3200
fix: make OutputRequirementExec execute transparently; pin test targe…
zhuqi-lucas 82f80ae
fix: clear analyzer rules in tests that observe an unoptimized physic…
zhuqi-lucas 51642ae
fix: pin target_partitions in window_topn / join_selection tests
zhuqi-lucas 2e3668a
fix: keep enforcement out of memory_limit scenarios that pin custom r…
zhuqi-lucas 0bf1852
fix: apply prettier formatting to optimizer_rule_reference.md
zhuqi-lucas f9cb358
docs: address review feedback on the analyzer-phase docs
zhuqi-lucas 64f1efb
Merge remote-tracking branch 'apache/main' into enforce-first-full
zhuqi-lucas 9d47f77
Merge branch 'main' into physical-analyzer-rule
zhuqi-lucas e0c199a
docs: add a phase diagram for the physical analyzer / optimizer split
zhuqi-lucas c6f93f6
Merge branch 'main' into physical-analyzer-rule
alamb 20c7daf
Run the rules that decide requirements before enforcing them
zhuqi-lucas b211ce1
Check the analyzer phase's contract where the phase ends
zhuqi-lucas 67d6f4a
Drop stale comments that still describe enforcement as an optimizer rule
zhuqi-lucas 9a39fe3
Document the physical analyzer phase in the 56.0.0 upgrade guide
zhuqi-lucas 267e70c
Keep EnsureRequirements as the single enforcement rule
zhuqi-lucas 64e8179
Merge branch 'physical-analyzer-rule' of github.com:zhuqi-lucas/arrow…
zhuqi-lucas bc862ce
Describe the enforcement phases as the two walks the code runs
zhuqi-lucas 86fde11
Reduce the analyzer split to the trait, the plumbing and today's rule…
zhuqi-lucas 79c63da
Merge branch 'main' into physical-analyzer-rule
zhuqi-lucas File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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}, | ||
|
|
@@ -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 | ||
| physical_analyzers: PhysicalAnalyzer, | ||
| /// Responsible for optimizing a physical execution plan | ||
| physical_optimizers: PhysicalOptimizer, | ||
| /// Responsible for planning `LogicalPlan`s, and `ExecutionPlan` | ||
|
|
@@ -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) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 |
||
| .field("physical_optimizers", &self.inner.physical_optimizers) | ||
| .field("table_functions", &self.inner.table_functions) | ||
| .field("scalar_functions", &self.inner.scalar_functions) | ||
|
|
@@ -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) | ||
| } | ||
|
|
@@ -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 | ||
|
|
@@ -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() | ||
|
|
@@ -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>>, | ||
|
|
@@ -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>>>, | ||
| } | ||
|
|
||
|
|
@@ -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, | ||
|
|
@@ -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, | ||
| } | ||
| } | ||
|
|
@@ -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), | ||
|
|
@@ -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, | ||
| } | ||
| } | ||
|
|
@@ -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, | ||
|
|
@@ -1689,6 +1745,7 @@ impl SessionStateBuilder { | |
| #[cfg(feature = "sql")] | ||
| type_planner, | ||
| optimizer, | ||
| physical_analyzers, | ||
| physical_optimizers, | ||
| query_planner, | ||
| catalog_list, | ||
|
|
@@ -1711,6 +1768,7 @@ impl SessionStateBuilder { | |
| statistics_registry, | ||
| analyzer_rules, | ||
| optimizer_rules, | ||
| physical_analyzer_rules, | ||
| physical_optimizer_rules, | ||
| } = self; | ||
|
|
||
|
|
@@ -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 {})), | ||
|
|
@@ -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; | ||
|
|
@@ -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 | ||
|
|
@@ -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 { | ||
|
|
@@ -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) | ||
|
|
||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
❤️