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]

Reply via email to