This is an automated email from the ASF dual-hosted git repository. gabotechs pushed a commit to branch gabotechs/make-query-planner-accept-dyn-session in repository https://gitbox.apache.org/repos/asf/datafusion.git
commit e262fbc35e22d4a82125c3c64dbbb2be32e04f00 Author: Gabriel <[email protected]> AuthorDate: Thu Jul 9 13:42:37 2026 +0200 Make query planner accept `dyn Session` --- datafusion/core/src/execution/context/mod.rs | 10 +-- datafusion/core/src/execution/session_state.rs | 2 +- datafusion/core/src/physical_planner.rs | 82 ++++++++++++++-------- .../core/tests/user_defined/user_defined_plan.rs | 7 +- 4 files changed, 62 insertions(+), 39 deletions(-) diff --git a/datafusion/core/src/execution/context/mod.rs b/datafusion/core/src/execution/context/mod.rs index 0ff3ab7d0e..6a3455b720 100644 --- a/datafusion/core/src/execution/context/mod.rs +++ b/datafusion/core/src/execution/context/mod.rs @@ -100,7 +100,7 @@ use datafusion_optimizer::analyzer::type_coercion::TypeCoercion; use datafusion_optimizer::simplify_expressions::ExprSimplifier; use datafusion_optimizer::{Analyzer, OptimizerContext}; use datafusion_optimizer::{AnalyzerRule, OptimizerRule}; -use datafusion_session::SessionStore; +use datafusion_session::{Session, SessionStore}; use async_trait::async_trait; use chrono::{DateTime, Utc}; @@ -2194,7 +2194,7 @@ pub trait QueryPlanner: Any + Debug { async fn create_physical_plan( &self, logical_plan: &LogicalPlan, - session_state: &SessionState, + session_state: &dyn Session, ) -> Result<Arc<dyn ExecutionPlan>>; } @@ -2841,7 +2841,7 @@ mod tests { async fn create_physical_plan( &self, _logical_plan: &LogicalPlan, - _session_state: &SessionState, + _session_state: &dyn Session, ) -> Result<Arc<dyn ExecutionPlan>> { not_impl_err!("query not supported") } @@ -2850,7 +2850,7 @@ mod tests { &self, _expr: &Expr, _input_dfschema: &DFSchema, - _session_state: &SessionState, + _session_state: &dyn Session, ) -> Result<Arc<dyn PhysicalExpr>> { unimplemented!() } @@ -2864,7 +2864,7 @@ mod tests { async fn create_physical_plan( &self, logical_plan: &LogicalPlan, - session_state: &SessionState, + session_state: &dyn Session, ) -> Result<Arc<dyn ExecutionPlan>> { let physical_planner = MyPhysicalPlanner {}; physical_planner diff --git a/datafusion/core/src/execution/session_state.rs b/datafusion/core/src/execution/session_state.rs index f1f5465212..7411c8cf58 100644 --- a/datafusion/core/src/execution/session_state.rs +++ b/datafusion/core/src/execution/session_state.rs @@ -2311,7 +2311,7 @@ impl QueryPlanner for DefaultQueryPlanner { async fn create_physical_plan( &self, logical_plan: &LogicalPlan, - session_state: &SessionState, + session_state: &dyn Session, ) -> datafusion_common::Result<Arc<dyn ExecutionPlan>> { let planner = DefaultPhysicalPlanner::default(); planner diff --git a/datafusion/core/src/physical_planner.rs b/datafusion/core/src/physical_planner.rs index b6d28e7b21..c17caef3f5 100644 --- a/datafusion/core/src/physical_planner.rs +++ b/datafusion/core/src/physical_planner.rs @@ -72,7 +72,7 @@ use datafusion_common::tree_node::{ }; use datafusion_common::{ DFSchema, DFSchemaRef, ScalarValue, exec_err, internal_datafusion_err, internal_err, - not_impl_err, plan_err, + not_impl_err, plan_datafusion_err, plan_err, }; use datafusion_common::{ TableReference, assert_eq_or_internal_err, assert_or_internal_err, @@ -112,6 +112,7 @@ use datafusion_physical_plan::unnest::ListUnnest; use async_trait::async_trait; use datafusion_physical_plan::async_func::{AsyncFuncExec, AsyncMapper}; +use datafusion_session::Session; use futures::{StreamExt, TryStreamExt}; use indexmap::IndexSet; use itertools::{Itertools, multiunzip}; @@ -126,7 +127,7 @@ pub trait PhysicalPlanner: Send + Sync { async fn create_physical_plan( &self, logical_plan: &LogicalPlan, - session_state: &SessionState, + session_state: &dyn Session, ) -> Result<Arc<dyn ExecutionPlan>>; /// Create a physical expression from a logical expression @@ -139,7 +140,7 @@ pub trait PhysicalPlanner: Send + Sync { &self, expr: &Expr, input_dfschema: &DFSchema, - session_state: &SessionState, + session_state: &dyn Session, ) -> Result<Arc<dyn PhysicalExpr>>; } @@ -162,7 +163,7 @@ pub trait ExtensionPlanner { node: &dyn UserDefinedLogicalNode, logical_inputs: &[&LogicalPlan], physical_inputs: &[Arc<dyn ExecutionPlan>], - session_state: &SessionState, + session_state: &dyn Session, ) -> Result<Option<Arc<dyn ExecutionPlan>>>; /// Create a physical plan for a [`LogicalPlan::TableScan`]. @@ -180,7 +181,7 @@ pub trait ExtensionPlanner { /// use std::sync::Arc; /// use datafusion::physical_plan::ExecutionPlan; /// use datafusion::logical_expr::TableScan; - /// use datafusion::execution::context::SessionState; + /// use datafusion_session::Session; /// use datafusion::error::Result; /// use datafusion_physical_planner::{ExtensionPlanner, PhysicalPlanner}; /// use async_trait::async_trait; @@ -201,7 +202,7 @@ pub trait ExtensionPlanner { /// _node: &dyn UserDefinedLogicalNode, /// _logical_inputs: &[&LogicalPlan], /// _physical_inputs: &[Arc<dyn ExecutionPlan>], - /// _session_state: &SessionState, + /// _session_state: &dyn Session, /// ) -> Result<Option<Arc<dyn ExecutionPlan>>> { /// Ok(None) /// } @@ -210,7 +211,7 @@ pub trait ExtensionPlanner { /// &self, /// _planner: &dyn PhysicalPlanner, /// scan: &TableScan, - /// _session_state: &SessionState, + /// _session_state: &dyn Session, /// ) -> Result<Option<Arc<dyn ExecutionPlan>>> { /// // Check if this is your custom table source /// if scan.source.is::<MyCustomTableSource>() { @@ -234,7 +235,7 @@ pub trait ExtensionPlanner { &self, _planner: &dyn PhysicalPlanner, _scan: &TableScan, - _session_state: &SessionState, + _session_state: &dyn Session, ) -> Result<Option<Arc<dyn ExecutionPlan>>> { Ok(None) } @@ -269,8 +270,17 @@ impl PhysicalPlanner for DefaultPhysicalPlanner { async fn create_physical_plan( &self, logical_plan: &LogicalPlan, - session_state: &SessionState, + session_state: &dyn Session, ) -> Result<Arc<dyn ExecutionPlan>> { + let session_state = session_state + .as_any() + .downcast_ref::<SessionState>() + .ok_or_else(|| { + plan_datafusion_err!( + "DefaultPhysicalPlanner requires a SessionState for physical optimization" + ) + })?; + if let Some(plan) = self .handle_explain_or_analyze(logical_plan, session_state) .await? @@ -294,7 +304,7 @@ impl PhysicalPlanner for DefaultPhysicalPlanner { &self, expr: &Expr, input_dfschema: &DFSchema, - session_state: &SessionState, + session_state: &dyn Session, ) -> Result<Arc<dyn PhysicalExpr>> { create_physical_expr(expr, input_dfschema, session_state.execution_props()) } @@ -435,7 +445,7 @@ impl DefaultPhysicalPlanner { fn create_initial_plan<'a>( &'a self, logical_plan: &'a LogicalPlan, - session_state: &'a SessionState, + session_state: &'a dyn Session, ) -> futures::future::BoxFuture<'a, Result<Arc<dyn ExecutionPlan>>> { Box::pin(async move { // When `enable_physical_uncorrelated_scalar_subquery` is disabled, the @@ -458,7 +468,11 @@ impl DefaultPhysicalPlanner { if links.is_empty() { return self - .create_initial_plan_inner(logical_plan, session_state) + .create_initial_plan_inner( + logical_plan, + session_state, + session_state.execution_props(), + ) .await; } @@ -472,13 +486,12 @@ impl DefaultPhysicalPlanner { // context rather than in `ExecutionProps`. It's here because // `create_physical_expr` only receives `&ExecutionProps`. let results = ScalarSubqueryResults::new(links.len()); - let mut owned = session_state.clone(); - owned.execution_props_mut().subquery_indexes = index_map; - owned.execution_props_mut().subquery_results = results.clone(); - let session_state = Cow::Owned(owned); + let mut props = session_state.execution_props().clone(); + props.subquery_indexes = index_map; + props.subquery_results = results.clone(); let plan = self - .create_initial_plan_inner(logical_plan, &session_state) + .create_initial_plan_inner(logical_plan, session_state, &props) .await?; Ok(Arc::new(ScalarSubqueryExec::new(plan, links, results))) }) @@ -489,7 +502,8 @@ impl DefaultPhysicalPlanner { async fn create_initial_plan_inner( &self, logical_plan: &LogicalPlan, - session_state: &SessionState, + session_state: &dyn Session, + props: &ExecutionProps, ) -> Result<Arc<dyn ExecutionPlan>> { // DFS the tree to flatten it into a Vec. // This will allow us to build the Physical Plan from the leaves up @@ -540,9 +554,9 @@ impl DefaultPhysicalPlanner { let max_concurrency = planning_concurrency.min(flat_tree_leaf_indices.len()); // Spawning tasks which will traverse leaf up to the root. - let tasks = flat_tree_leaf_indices - .into_iter() - .map(|index| self.task_helper(index, Arc::clone(&flat_tree), session_state)); + let tasks = flat_tree_leaf_indices.into_iter().map(|index| { + self.task_helper(index, Arc::clone(&flat_tree), session_state, props) + }); let mut outputs = futures::stream::iter(tasks) .buffer_unordered(max_concurrency) .try_collect::<Vec<_>>() @@ -569,7 +583,8 @@ impl DefaultPhysicalPlanner { &'a self, leaf_starter_index: usize, flat_tree: Arc<Vec<LogicalNode<'a>>>, - session_state: &'a SessionState, + session_state: &'a dyn Session, + props: &ExecutionProps, ) -> Result<Option<Arc<dyn ExecutionPlan>>> { // We always start with a leaf, so can ignore status and pass empty children let mut node = flat_tree.get(leaf_starter_index).ok_or_else(|| { @@ -582,6 +597,7 @@ impl DefaultPhysicalPlanner { node.node, session_state, ChildrenContainer::None, + props, ) .await?; let mut current_index = leaf_starter_index; @@ -599,6 +615,7 @@ impl DefaultPhysicalPlanner { node.node, session_state, ChildrenContainer::One(plan), + props, ) .await?; } @@ -634,7 +651,12 @@ impl DefaultPhysicalPlanner { let children = children.into_iter().map(|epc| epc.plan).collect(); let children = ChildrenContainer::Multiple(children); plan = self - .map_logical_node_to_physical(node.node, session_state, children) + .map_logical_node_to_physical( + node.node, + session_state, + children, + props, + ) .await?; } } @@ -648,10 +670,10 @@ impl DefaultPhysicalPlanner { async fn map_logical_node_to_physical( &self, node: &LogicalPlan, - session_state: &SessionState, + session_state: &dyn Session, children: ChildrenContainer, + execution_props: &ExecutionProps, ) -> Result<Arc<dyn ExecutionPlan>> { - let execution_props = session_state.execution_props(); let exec_node: Arc<dyn ExecutionPlan> = match node { // Leaves (no children) LogicalPlan::TableScan(scan) => { @@ -2928,7 +2950,7 @@ impl DefaultPhysicalPlanner { async fn plan_scalar_subqueries( &self, subqueries: Vec<Subquery>, - session_state: &SessionState, + session_state: &dyn Session, ) -> Result<(Vec<ScalarSubqueryLink>, DFHashMap<Subquery, SubqueryIndex>)> { let mut links = Vec::with_capacity(subqueries.len()); let mut index_map = DFHashMap::with_capacity(subqueries.len()); @@ -4395,7 +4417,7 @@ mod tests { _node: &dyn UserDefinedLogicalNode, _logical_inputs: &[&LogicalPlan], _physical_inputs: &[Arc<dyn ExecutionPlan>], - _session_state: &SessionState, + _session_state: &dyn Session, ) -> Result<Option<Arc<dyn ExecutionPlan>>> { internal_err!("BOOM") } @@ -4549,7 +4571,7 @@ mod tests { _node: &dyn UserDefinedLogicalNode, _logical_inputs: &[&LogicalPlan], _physical_inputs: &[Arc<dyn ExecutionPlan>], - _session_state: &SessionState, + _session_state: &dyn Session, ) -> Result<Option<Arc<dyn ExecutionPlan>>> { Ok(Some(Arc::new(NoOpExecutionPlan::new(SchemaRef::new( Schema::new(vec![Field::new("b", DataType::Int32, false)]), @@ -5183,7 +5205,7 @@ digraph { _node: &dyn UserDefinedLogicalNode, _logical_inputs: &[&LogicalPlan], _physical_inputs: &[Arc<dyn ExecutionPlan>], - _session_state: &SessionState, + _session_state: &dyn Session, ) -> Result<Option<Arc<dyn ExecutionPlan>>> { Ok(None) } @@ -5192,7 +5214,7 @@ digraph { &self, _planner: &dyn PhysicalPlanner, scan: &TableScan, - _session_state: &SessionState, + _session_state: &dyn Session, ) -> Result<Option<Arc<dyn ExecutionPlan>>> { if scan.source.is::<MockTableSource>() { Ok(Some(Arc::new(EmptyExec::new(Arc::clone( diff --git a/datafusion/core/tests/user_defined/user_defined_plan.rs b/datafusion/core/tests/user_defined/user_defined_plan.rs index e8ff6758cc..81655471ea 100644 --- a/datafusion/core/tests/user_defined/user_defined_plan.rs +++ b/datafusion/core/tests/user_defined/user_defined_plan.rs @@ -69,11 +69,12 @@ use arrow::{ }; use datafusion::execution::session_state::SessionStateBuilder; use datafusion::{ + catalog::Session, common::cast::as_int64_array, common::{DFSchemaRef, arrow_datafusion_err}, error::{DataFusionError, Result}, execution::{ - context::{QueryPlanner, SessionState, TaskContext}, + context::{QueryPlanner, TaskContext}, runtime_env::RuntimeEnv, }, logical_expr::{ @@ -466,7 +467,7 @@ impl QueryPlanner for TopKQueryPlanner { async fn create_physical_plan( &self, logical_plan: &LogicalPlan, - session_state: &SessionState, + session_state: &dyn Session, ) -> Result<Arc<dyn ExecutionPlan>> { // Teach the default physical planner how to plan TopK nodes. let physical_planner = @@ -629,7 +630,7 @@ impl ExtensionPlanner for TopKPlanner { node: &dyn UserDefinedLogicalNode, logical_inputs: &[&LogicalPlan], physical_inputs: &[Arc<dyn ExecutionPlan>], - _session_state: &SessionState, + _session_state: &dyn Session, ) -> Result<Option<Arc<dyn ExecutionPlan>>> { Ok( if let Some(topk_node) = node.as_any().downcast_ref::<TopKPlanNode>() { --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
