Skip to content

Commit 3179bd9

Browse files
committed
perf: share one StatisticsContext across all physical optimizer rules
Create one StatisticsContext per optimize_physical_plan call and expose it through the new PhysicalOptimizerContext::statistics_context, so statistics computed by one rule are reused by later rules. JoinSelection, EnsureRequirements, AggregateStatistics and LimitPushdown use it. StatisticsContext stores its cache in a parking_lot::Mutex so it is Send + Sync, as PhysicalOptimizerContext requires.
1 parent db83fcc commit 3179bd9

9 files changed

Lines changed: 282 additions & 120 deletions

File tree

‎datafusion/core/src/physical_planner.rs‎

Lines changed: 84 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -106,6 +106,7 @@ use datafusion_physical_plan::joins::PiecewiseMergeJoinExec;
106106
use datafusion_physical_plan::placeholder_row::PlaceholderRowExec;
107107
use datafusion_physical_plan::recursive_query::RecursiveQueryExec;
108108
use datafusion_physical_plan::scalar_subquery::{ScalarSubqueryExec, ScalarSubqueryLink};
109+
use datafusion_physical_plan::statistics::StatisticsContext;
109110
use datafusion_physical_plan::unnest::ListUnnest;
110111
use datafusion_session::{PhysicalOptimizerContext, PhysicalOptimizerRule, Session};
111112

@@ -123,6 +124,20 @@ pub use datafusion_session::{ExtensionPlanner, PhysicalPlanner};
123124

124125
struct SessionOptimizerContext<'a> {
125126
session: &'a dyn Session,
127+
statistics_context: StatisticsContext,
128+
}
129+
130+
impl<'a> SessionOptimizerContext<'a> {
131+
fn new(session: &'a dyn Session) -> Self {
132+
let statistics_context = match session.statistics_registry() {
133+
Some(registry) => StatisticsContext::new_with_registry(registry.clone()),
134+
None => StatisticsContext::new(),
135+
};
136+
Self {
137+
session,
138+
statistics_context,
139+
}
140+
}
126141
}
127142

128143
impl PhysicalOptimizerContext for SessionOptimizerContext<'_> {
@@ -135,6 +150,10 @@ impl PhysicalOptimizerContext for SessionOptimizerContext<'_> {
135150
) -> Option<&datafusion_physical_plan::operator_statistics::StatisticsRegistry> {
136151
self.session.statistics_registry()
137152
}
153+
154+
fn statistics_context(&self) -> Option<&StatisticsContext> {
155+
Some(&self.statistics_context)
156+
}
138157
}
139158

140159
/// Default single node physical query planner that converts a
@@ -3119,9 +3138,7 @@ impl DefaultPhysicalPlanner {
31193138
InvariantChecker(InvariantLevel::Always).check(&plan)?;
31203139

31213140
let mut new_plan = Arc::clone(&plan);
3122-
let optimizer_context = SessionOptimizerContext {
3123-
session: session_state,
3124-
};
3141+
let optimizer_context = SessionOptimizerContext::new(session_state);
31253142
for optimizer in optimizers {
31263143
let before_schema = new_plan.schema();
31273144
new_plan = optimizer
@@ -3532,6 +3549,7 @@ mod tests {
35323549
use arrow::datatypes::{DataType, Field, Int32Type};
35333550
use arrow_schema::{FieldRef, SchemaRef};
35343551
use datafusion_catalog::CatalogProviderList;
3552+
use datafusion_common::Statistics;
35353553
use datafusion_common::config::{ConfigOptions, TableOptions};
35363554
use datafusion_common::{
35373555
DFSchemaRef, ScalarValue, SplitPoint, TableReference, ToDFSchema as _,
@@ -3554,8 +3572,10 @@ mod tests {
35543572
use datafusion_functions_aggregate::expr_fn::sum;
35553573
use datafusion_physical_expr::EquivalenceProperties;
35563574
use datafusion_physical_plan::execution_plan::{Boundedness, EmissionType};
3575+
use datafusion_physical_plan::statistics::StatisticsArgs;
35573576
use datafusion_physical_plan::{ChildrenPropertiesMode, ReplaceChildrenOptions};
35583577
use datafusion_session::QueryPlanner;
3578+
use parking_lot::Mutex as SyncMutex;
35593579

35603580
#[derive(Debug)]
35613581
struct ContextCheckingRule {
@@ -3590,6 +3610,44 @@ mod tests {
35903610
}
35913611
}
35923612

3613+
/// Records the root statistics computed with the shared statistics context
3614+
#[derive(Debug)]
3615+
struct StatisticsRecordingRule {
3616+
recorded: Arc<SyncMutex<Vec<Arc<Statistics>>>>,
3617+
}
3618+
3619+
impl PhysicalOptimizerRule for StatisticsRecordingRule {
3620+
fn optimize(
3621+
&self,
3622+
plan: Arc<dyn ExecutionPlan>,
3623+
_config: &ConfigOptions,
3624+
) -> Result<Arc<dyn ExecutionPlan>> {
3625+
Ok(plan)
3626+
}
3627+
3628+
fn optimize_with_context(
3629+
&self,
3630+
plan: Arc<dyn ExecutionPlan>,
3631+
context: &dyn PhysicalOptimizerContext,
3632+
) -> Result<Arc<dyn ExecutionPlan>> {
3633+
let statistics_context = context
3634+
.statistics_context()
3635+
.expect("the planner shares a statistics context");
3636+
let statistics =
3637+
statistics_context.compute_arc(&plan, &StatisticsArgs::new())?;
3638+
self.recorded.lock().push(statistics);
3639+
Ok(plan)
3640+
}
3641+
3642+
fn name(&self) -> &str {
3643+
"statistics_recording_rule"
3644+
}
3645+
3646+
fn schema_check(&self) -> bool {
3647+
true
3648+
}
3649+
}
3650+
35933651
#[derive(Debug)]
35943652
struct TestQueryPlanner {
35953653
invoked: Arc<AtomicBool>,
@@ -3839,6 +3897,29 @@ mod tests {
38393897
Ok(())
38403898
}
38413899

3900+
#[tokio::test]
3901+
async fn optimizer_rules_share_statistics_context() -> Result<()> {
3902+
let recorded = Arc::new(SyncMutex::new(vec![]));
3903+
let rule = || -> Arc<dyn PhysicalOptimizerRule + Send + Sync> {
3904+
Arc::new(StatisticsRecordingRule {
3905+
recorded: Arc::clone(&recorded),
3906+
})
3907+
};
3908+
let session_state = SessionStateBuilder::new()
3909+
.with_default_features()
3910+
.with_physical_optimizer_rules(vec![rule(), rule()])
3911+
.build();
3912+
3913+
let logical_plan = LogicalPlanBuilder::empty(false).build()?;
3914+
session_state.create_physical_plan(&logical_plan).await?;
3915+
3916+
// The second rule reads the statistics the first rule cached
3917+
let recorded = recorded.lock();
3918+
assert_eq!(recorded.len(), 2);
3919+
assert!(Arc::ptr_eq(&recorded[0], &recorded[1]));
3920+
Ok(())
3921+
}
3922+
38423923
async fn aggregate_explain(logical_plan: &LogicalPlan) -> Result<String> {
38433924
let physical_plan = plan(logical_plan).await?;
38443925
Ok(displayable(physical_plan.as_ref()).indent(true).to_string())

‎datafusion/physical-optimizer/src/aggregate_statistics.rs‎

Lines changed: 23 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -25,14 +25,17 @@ use datafusion_physical_plan::aggregates::{
2525
};
2626
use datafusion_physical_plan::placeholder_row::PlaceholderRowExec;
2727
use datafusion_physical_plan::projection::{ProjectionExec, ProjectionExpr};
28-
use datafusion_physical_plan::statistics::{StatisticsArgs, StatisticsContext};
28+
use datafusion_physical_plan::statistics::StatisticsArgs;
2929
use datafusion_physical_plan::udaf::{
3030
AggregateFunctionExpr, StatisticsArgs as PlanStatisticsArgs,
3131
};
3232
use datafusion_physical_plan::{ExecutionPlan, expressions};
3333
use std::sync::Arc;
3434

3535
use crate::PhysicalOptimizerRule;
36+
use crate::optimizer::{
37+
ConfigOnlyContext, PhysicalOptimizerContext, with_statistics_context,
38+
};
3639

3740
/// Optimizer that uses available statistics for aggregate functions
3841
#[derive(Default, Debug)]
@@ -46,20 +49,27 @@ impl AggregateStatistics {
4649
}
4750

4851
impl PhysicalOptimizerRule for AggregateStatistics {
49-
#[cfg_attr(feature = "recursive_protection", recursive::recursive)]
50-
#[expect(clippy::allow_attributes)] // See https://github.com/apache/datafusion/issues/18881#issuecomment-3621545670
51-
#[allow(clippy::only_used_in_recursion)] // See https://github.com/rust-lang/rust-clippy/issues/14566
5252
fn optimize(
5353
&self,
5454
plan: Arc<dyn ExecutionPlan>,
5555
config: &ConfigOptions,
56+
) -> Result<Arc<dyn ExecutionPlan>> {
57+
self.optimize_with_context(plan, &ConfigOnlyContext::new(config))
58+
}
59+
60+
#[cfg_attr(feature = "recursive_protection", recursive::recursive)]
61+
fn optimize_with_context(
62+
&self,
63+
plan: Arc<dyn ExecutionPlan>,
64+
context: &dyn PhysicalOptimizerContext,
5665
) -> Result<Arc<dyn ExecutionPlan>> {
5766
if let Some(partial_agg_exec) = take_optimizable(&plan) {
5867
let partial_agg_exec = partial_agg_exec
5968
.downcast_ref::<AggregateExec>()
6069
.expect("take_optimizable() ensures that this is a AggregateExec");
61-
let stats = StatisticsContext::new()
62-
.compute(partial_agg_exec.input().as_ref(), &StatisticsArgs::new())?;
70+
let stats = with_statistics_context(context, |stats_ctx| {
71+
stats_ctx.compute_arc(partial_agg_exec.input(), &StatisticsArgs::new())
72+
})?;
6373
let mut projections = vec![];
6474
for expr in partial_agg_exec.aggr_expr() {
6575
let field = expr.field();
@@ -92,13 +102,17 @@ impl PhysicalOptimizerRule for AggregateStatistics {
92102
)?))
93103
} else {
94104
plan.map_children(|child| {
95-
self.optimize(child, config).map(Transformed::yes)
105+
self.optimize_with_context(child, context)
106+
.map(Transformed::yes)
96107
})
97108
.data()
98109
}
99110
} else {
100-
plan.map_children(|child| self.optimize(child, config).map(Transformed::yes))
101-
.data()
111+
plan.map_children(|child| {
112+
self.optimize_with_context(child, context)
113+
.map(Transformed::yes)
114+
})
115+
.data()
102116
}
103117
}
104118

‎datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs‎

Lines changed: 48 additions & 49 deletions
Original file line numberDiff line numberDiff line change
@@ -1063,10 +1063,8 @@ struct PlanSize {
10631063
}
10641064

10651065
impl PlanSize {
1066-
fn from_plan(plan: &dyn ExecutionPlan) -> Self {
1067-
let stats = StatisticsContext::new()
1068-
.compute(plan, &StatisticsArgs::new())
1069-
.ok();
1066+
fn from_plan(plan: &Arc<dyn ExecutionPlan>, stats_ctx: &StatisticsContext) -> Self {
1067+
let stats = stats_ctx.compute_arc(plan, &StatisticsArgs::new()).ok();
10701068
Self {
10711069
byte_size: stats
10721070
.as_ref()
@@ -1124,6 +1122,7 @@ fn enforce_distribution_relationships(
11241122
input_distributions: &InputDistributionRequirements,
11251123
children: &mut [DistributionChildState],
11261124
target_partitions: usize,
1125+
stats_ctx: &StatisticsContext,
11271126
) -> Result<()> {
11281127
let mut repartitioned_for_relationship = vec![false; children.len()];
11291128

@@ -1162,53 +1161,52 @@ fn enforce_distribution_relationships(
11621161
}
11631162
}
11641163

1165-
let best_satisfied_child: Option<(usize, Partitioning)> = match satisfied_children
1166-
.len()
1167-
{
1168-
0 => None,
1169-
1 => satisfied_children
1170-
.into_iter()
1171-
.next()
1172-
.map(|(i, p, _)| (i, p)),
1173-
_ => {
1174-
// Prefer native partitioned children over newly repartitioned ones
1175-
let native_children: Vec<_> = satisfied_children
1176-
.iter()
1177-
.filter(|(_, _, is_native)| *is_native)
1178-
.collect();
1179-
if native_children.len() == 1 {
1180-
let (i, p, _) = native_children[0];
1181-
Some((*i, p.clone()))
1182-
} else {
1183-
let pool = if !native_children.is_empty() {
1184-
native_children
1185-
} else {
1186-
satisfied_children.iter().collect()
1187-
};
1188-
let candidates: Vec<_> = pool
1189-
.into_iter()
1190-
.map(|(idx, part, _)| {
1191-
let size =
1192-
PlanSize::from_plan(children[*idx].context.plan.as_ref());
1193-
(size, *idx, part.clone())
1194-
})
1195-
.collect();
1196-
1197-
// Prefer a unique, strictly larger winner (`size_a > size_b`).
1198-
// Otherwise, fall back to standard distribution rather
1199-
// than choosing an arbitrary reference.
1200-
candidates
1164+
let best_satisfied_child: Option<(usize, Partitioning)> =
1165+
match satisfied_children.len() {
1166+
0 => None,
1167+
1 => satisfied_children
1168+
.into_iter()
1169+
.next()
1170+
.map(|(i, p, _)| (i, p)),
1171+
_ => {
1172+
// Prefer native partitioned children over newly repartitioned ones
1173+
let native_children: Vec<_> = satisfied_children
12011174
.iter()
1202-
.find(|(size_a, idx_a, _)| {
1203-
size_a.is_known()
1204-
&& candidates.iter().all(|(size_b, idx_b, _)| {
1205-
idx_a == idx_b || size_a > size_b
1206-
})
1207-
})
1208-
.map(|(_, idx, part)| (*idx, part.clone()))
1175+
.filter(|(_, _, is_native)| *is_native)
1176+
.collect();
1177+
if native_children.len() == 1 {
1178+
let (i, p, _) = native_children[0];
1179+
Some((*i, p.clone()))
1180+
} else {
1181+
let pool = if !native_children.is_empty() {
1182+
native_children
1183+
} else {
1184+
satisfied_children.iter().collect()
1185+
};
1186+
let candidates: Vec<_> = pool
1187+
.into_iter()
1188+
.map(|(idx, part, _)| {
1189+
let plan = &children[*idx].context.plan;
1190+
let size = PlanSize::from_plan(plan, stats_ctx);
1191+
(size, *idx, part.clone())
1192+
})
1193+
.collect();
1194+
1195+
// Prefer a unique, strictly larger winner (`size_a > size_b`).
1196+
// Otherwise, fall back to standard distribution rather
1197+
// than choosing an arbitrary reference.
1198+
candidates
1199+
.iter()
1200+
.find(|(size_a, idx_a, _)| {
1201+
size_a.is_known()
1202+
&& candidates.iter().all(|(size_b, idx_b, _)| {
1203+
idx_a == idx_b || size_a > size_b
1204+
})
1205+
})
1206+
.map(|(_, idx, part)| (*idx, part.clone()))
1207+
}
12091208
}
1210-
}
1211-
};
1209+
};
12121210

12131211
// Validate that best_satisfied_child can be adapted across all unsatisfied children in the group
12141212
let best_satisfied_child =
@@ -1666,6 +1664,7 @@ pub fn ensure_distribution_with_stats(
16661664
&input_distributions,
16671665
&mut children,
16681666
target_partitions,
1667+
stats_ctx,
16691668
)?;
16701669

16711670
let children = children

‎datafusion/physical-optimizer/src/ensure_requirements/mod.rs‎

Lines changed: 10 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -150,13 +150,14 @@ pub mod enforce_sorting;
150150
use std::sync::Arc;
151151

152152
use crate::PhysicalOptimizerRule;
153-
use crate::optimizer::{ConfigOnlyContext, PhysicalOptimizerContext};
153+
use crate::optimizer::{
154+
ConfigOnlyContext, PhysicalOptimizerContext, with_statistics_context,
155+
};
154156

155157
use datafusion_common::Result;
156158
use datafusion_common::config::ConfigOptions;
157159
use datafusion_common::tree_node::{Transformed, TransformedResult, TreeNode};
158160
use datafusion_physical_plan::ExecutionPlan;
159-
use datafusion_physical_plan::statistics::StatisticsContext;
160161

161162
/// Optimizer rule that enforces both distribution and sorting requirements.
162163
///
@@ -221,17 +222,13 @@ impl PhysicalOptimizerRule for EnsureRequirements {
221222

222223
// Step 2a: Distribution enforcement (bottom-up)
223224
let dist_ctx = DistributionContext::new_default(plan);
224-
// Share one statistics context across the whole distribution pass so each
225-
// subtree's statistics are computed once instead of once per ancestor.
226-
// Build it from the session's statistics registry so registered providers
227-
// are consulted (an empty registry, the default, is unchanged behavior).
228-
let stats_ctx = match context.statistics_registry() {
229-
Some(registry) => StatisticsContext::new_with_registry(registry.clone()),
230-
None => StatisticsContext::new(),
231-
};
232-
let dist_ctx = dist_ctx
233-
.transform_up(|ctx| ensure_distribution_with_stats(ctx, config, &stats_ctx))
234-
.data()?;
225+
let dist_ctx = with_statistics_context(context, |stats_ctx| {
226+
dist_ctx
227+
.transform_up(|ctx| {
228+
ensure_distribution_with_stats(ctx, config, stats_ctx)
229+
})
230+
.data()
231+
})?;
235232

236233
// Step 2b: Sorting enforcement (bottom-up) — runs on distribution-fixed plan
237234
let sort_ctx = PlanWithCorrespondingSort::new_default(dist_ctx.plan);

0 commit comments

Comments
 (0)