Repository navigation
perf: share one StatisticsContext across all physical optimizer rules - #26094
asolimando wants to merge 8 commits into
Conversation
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.
|
@kosiew @zhuqi-lucas, could you run I figure you'd be interested in this PR as it builds on #25929, and it's the query-lifetime follow-up of what #25098 did for EnsureDistribution. |
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #26094 +/- ##
=========================================
Coverage 82.72% 82.72%
=========================================
Files 1147 1147
Lines 448179 449304 +1125
Branches 448179 449304 +1125
=========================================
+ Hits 370754 371701 +947
- Misses 54895 54947 +52
- Partials 22530 22656 +126 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
run benchmark sql_planner |
@asolimando triggered now |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing asolimando/query-lifetime-stats-cache (3179bd9) to db83fcc (merge-base) diff Run configurationrun benchmark sql_plannerResults will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing asolimando/query-lifetime-stats-cache (3179bd9) to db83fcc (merge-base) diff Run configurationrun benchmark sql_plannerCPU Details (lscpu)Details
Resource Usagesql_planner — base (merge-base)
sql_planner — branch
File an issue against this benchmark runner |
|
Thanks @zhuqi-lucas! My understanding of the benchmark run:
The +5% on The gain is little smaller than the sf1 numbers I had locally (6%, as reported in the description) because the benchmark tables are empty, so each statistics computation is cheap. If you think this is interesting I can do another self-review pass on the PR and remove it from draft, wdyt? |
|
Thanks for the detailed breakdown, your reading matches mine. This looks valuable, please go ahead with the self-review and take it out of draft. |
|
Thank you for opening this pull request! Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch). Details |
|
@zhuqi-lucas thanks a lot for your feedback, after self-review I have pushed some test refactoring, doc improvements and improved API ergonomics (added I noticed the coverage warning, I could fix that easily but I don't see anything that genuinely needs more coverage, happy to add more tests if you see fit. OT: is there a process to get into the allow-list to run benchmarks? What are the criteria for eligibility? |
| /// [`Statistics`] only. | ||
| pub struct StatisticsContext { | ||
| cache: Rc<RefCell<StatsCache>>, | ||
| cache: Mutex<StatsCache>, |
There was a problem hiding this comment.
Note for reviewers: the Mutex is needed because PhysicalOptimizerContext: Send + Sync, and the shared context is held inside it and returned as &StatisticsContext. RefCell is not Sync (and the Rc was never cloned, so it was dropped rather than replaced).
The optimizer uses the context from one thread, so the lock is never contended, and it is briefly held (single map lookup/insert), never across the recursive walk. The compute_statistics benchmark shows no measurable difference from RefCell (all cases within ±4% of main, in both directions).
Alternatives I considered: a RefCell restricted to the creating thread (needs unsafe impl Sync), a thread-local cache (hidden global state), or removing Sync from PhysicalOptimizerContext (a breaking change). None of them seemed worth it for a lock that costs nothing measurable.
| // Otherwise, fall back to standard distribution rather | ||
| // than choosing an arbitrary reference. | ||
| candidates | ||
| let best_satisfied_child: Option<(usize, Partitioning)> = |
There was a problem hiding this comment.
Note to reviewers: re-indented by rustfmt, the real change is passing stats_ctx to PlanSize::from_plan, but the diff can be confusing.
zhuqi-lucas
left a comment
There was a problem hiding this comment.
Nice work — this is the piece #25098 and #25929 were building toward, and the numbers are convincing.
I checked the two things that worried me most and both hold up:
- No lock held across the recursive walk. Every
borrow()→lock()site is a single lookup or insert whose guard dies at the end of the statement, so the non-reentrantparking_lot::Mutexcan't self-deadlock. - Pointer keys stay safe over the much longer sharing window.
store_cache_entryis the only insertion point and always clones the owner intoCacheEntry::_plan, so an address can't be recycled while its entry lives. That invariant is what makes widening the scope from one rule to the whole run sound — worth keeping in mind for anyone who later adds a second insertion path.
Also worth calling out, since it's easy to miss in the diff: the compute(x.as_ref()) → compute_arc(x) switches are not cosmetic. Borrowed roots are deliberately not memoized, so those lines are what actually lets the root of each traversal land in the shared cache.
Three small notes inline, none blocking.
… helper, ConfigOnlyContext invariant
|
run benchmark sql_planner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing asolimando/query-lifetime-stats-cache (63aa427) to 3ed377a (merge-base) diff Run configurationrun benchmark sql_plannerResults will be posted here when complete File an issue against this benchmark runner |
| let optimizer_context = SessionOptimizerContext { | ||
| session: session_state, | ||
| }; | ||
| let optimizer_context = SessionOptimizerContext::new(session_state); |
There was a problem hiding this comment.
One more thing this widens, worth stating somewhere: the cache is now exposed to in-place mutation for the whole planning run, not just one rule.
A pointer-keyed hit is only valid because a plan node never changes behind its Arc. That held trivially when each rule had its own context — anything a rule did produced new nodes. Now an entry computed by the first rule is still served to the last one, so a node that mutated its own statistics through interior mutability without changing identity would be read stale.
Nothing does that today (planning-time rules all rebuild nodes, and dynamic filters are updated at execution time, after this context is dropped), so this isn't a bug — but it's an invariant the design now leans on much harder. A line on StatisticsContext saying cached nodes must be immutable for the context's lifetime would make it checkable in review.
There was a problem hiding this comment.
Agreed, it is an invariant the design now relies on much more. Added a paragraph in f7587af to the StatisticsContext docs: a plan node must not change its statistics in place while a context holds it, otherwise the cache returns stale values; optimizer rules satisfy this because they replace nodes instead of changing them.
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing asolimando/query-lifetime-stats-cache (63aa427) to 3ed377a (merge-base) diff Run configurationrun benchmark sql_plannerCPU Details (lscpu)Details
Resource Usagesql_planner — base (merge-base)
sql_planner — branch
File an issue against this benchmark runner |
|
Thanks @zhuqi-lucas for re-running the benchmark! Here is how I read the second run:
On TPC-DS we save less time than in the first run (13 ms instead of 25 ms). I checked the commits merged to WDYT? |
Which issue does this PR close?
Rationale for this change
Each physical optimizer rule that reads statistics creates its own
StatisticsContext, so the statistics of the same plan nodes are computed again in every rule.JoinSelectioncreates a new context for everyget_statscall, so it has no caching at all, even within its own pass.#25098 showed the gain of sharing one context inside
EnsureDistribution. #25929 made cache entries keep the plan node they were computed for, so a context can now be shared safely across plan rewrites. This PR uses that to share oneStatisticsContextacross all the rules of oneoptimize_physical_plancall.On TPC-DS (98 queries) and TPC-H (21 queries), sf1 Parquet, planned with a local harness (not part of this PR):
mainThe largest gain is TPC-DS q64: cache misses go from 3,785 to 531 and planning time drops by about 21%.
What changes are included in this PR?
StatisticsContextstores its cache in aparking_lot::Mutexinstead ofRc<RefCell>, so it isSend + Sync. This is needed becausePhysicalOptimizerContext: Send + Sync. The lock is held only for single map lookups and inserts, never across the recursive walk. Thecompute_statisticsbenchmark shows no difference frommain(all cases within ±4%, in both directions).PhysicalOptimizerContext::statistics_context(), which returnsNoneby default.DefaultPhysicalPlannercreates one context peroptimize_physical_plancall, built from the session's statistics registry, and returns it to every rule.ConfigOnlyContextalso owns aStatisticsContext, so a rule called throughoptimize()shares one cache for its whole pass.JoinSelection,EnsureRequirements(includingPlanSize::from_plan),AggregateStatisticsandLimitPushdownuse the shared context. When a context does not share one (for example, a context received through FFI), the rules create a new context from the statistics registry, as they did before.PhysicalOptimizerContext::compute_statistics, which custom rules can call on the context they already receive. Without it, a custom rule would have to repeat the fallback itself (use the shared context, otherwise build one fromstatistics_registry()), and the obvious shortcut,StatisticsContext::new(), does not consult the registered statistics providers. With it, custom rules behave exactly like the built-in ones.with_statistics_context. Rules that walk the plan (EnsureRequirements,LimitPushdown) use it to get oneStatisticsContextfor the whole traversal, so a fallback context still caches within the pass.pushdown_limit_helperis deprecated in favour of the newpushdown_limit_helper_with_stats, which takes a&StatisticsContext. This is the same pattern asensure_distribution_with_stats.optimize_physical_planreturns, including nodes that a later rule replaces. Only the rules that read statistics add entries, and the memory is released at the end of physical planning.What is the testing strategy for this PR?
optimizer_rules_share_statistics_contexttest inphysical_planner.rs: two rules compute the root statistics through the shared context, and the second one gets theArccached by the first.Are there any user-facing changes?
No breaking changes.
PhysicalOptimizerContext::statistics_context()andPhysicalOptimizerContext::compute_statistics().datafusion_session::with_statistics_contextandpushdown_limit_helper_with_stats;pushdown_limit_helperis deprecated. All are described in the 56.0.0 upgrade guide.StatisticsContextis nowSend + Sync.AggregateStatistics,LimitPushdownandPlanSize::from_plannow consult the session's statistics providers, asJoinSelectionandEnsureRequirementsalready do. Without registered providers (the default), their results are unchanged.Disclaimer: I used AI to assist in the code generation, I have manually reviewed the output and it matches my intention and understanding.