Repository navigation
[X-3748] ensure_distribution: reuse one StatisticsContext across the pass - #78
Merged
Merged
Conversation
zhuqi-lucas
force-pushed
the
qizhu/x-3748
branch
from
September 9, 2026 06:57
c65be58 to
244d2c2
Compare
…pass Syncs this backport with what upstream actually merged in apache#25098, which grew substantially during review. ensure_distribution now takes its StatisticsContext from the caller, so one context and its memoization cache are shared across a whole bottom-up traversal instead of a fresh one per child, which recomputed each shared subtree's statistics once per ancestor. The old two-argument form stays as a deprecated wrapper. StatsCache is keyed by raw plan-node pointers, so the caller resets it after any node whose plan pointer actually changed: a rewrite can free a cached node and a later allocation could reuse its address. Three pieces of the upstream change are deliberately left out, since each depends on machinery this branch does not have and none of them is what the change is for: - The statistics registry threading. StatisticsContext here has no registry field and compute never consults providers, so the context is built empty. The memoization is unaffected. - The EnsureRequirements::optimize_with_context override, whose only purpose upstream was to reach that registry. - Phase 0's InterleaveExec normalization, which is context in the upstream diff rather than part of this change. Upstream's ensure_distribution_uses_context_statistics_registry test is dropped for the same reason: it exercises the registry path above. The fetch tests this branch carries from #75 are kept over the upstream versions they conflicted with; they are unrelated to this change. Verified: physical_optimizer::enforce_distribution passes 82/82, clippy and fmt clean. physical_optimizer::sanity_checker has 9 failures, which reproduce on a clean branch-55 checkout and are not from this change.
zhuqi-lucas
force-pushed
the
qizhu/x-3748
branch
from
September 17, 2026 03:42
244d2c2 to
e684ce4
Compare
zhuqi-lucas
marked this pull request as ready for review
September 17, 2026 03:43
There was a problem hiding this comment.
🟢 Approval recommended
The cache reuse and invalidation logic are consistent with the pointer-keyed cache, and regression coverage validates the intended optimization.
Pull request overview
This PR shares one StatisticsContext across the distribution-enforcement traversal, preserving memoized statistics and resetting safely after plan rewrites.
Changes:
- Added
ensure_distribution_with_statswith caller-provided statistics context. - Updated
EnsureRequirementsto reuse and invalidate the cache based on plan pointer changes. - Added regression coverage measuring statistics recomputation.
File summaries
| File | Description |
|---|---|
datafusion/physical-optimizer/src/ensure_requirements/mod.rs |
Shares the statistics context during distribution enforcement. |
datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs |
Threads the context through repartition decisions and adds a compatibility wrapper. |
datafusion/core/tests/physical_optimizer/enforce_distribution.rs |
Adds cache-sharing regression coverage. |
Review details
- Files reviewed: 3/3 changed files
- Comments generated: 0
- Review effort level: Lite
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
MassivePizza
pushed a commit
that referenced
this pull request
Sep 30, 2026
…pass (#78) Syncs this backport with what upstream actually merged in apache#25098, which grew substantially during review. ensure_distribution now takes its StatisticsContext from the caller, so one context and its memoization cache are shared across a whole bottom-up traversal instead of a fresh one per child, which recomputed each shared subtree's statistics once per ancestor. The old two-argument form stays as a deprecated wrapper. StatsCache is keyed by raw plan-node pointers, so the caller resets it after any node whose plan pointer actually changed: a rewrite can free a cached node and a later allocation could reuse its address. Three pieces of the upstream change are deliberately left out, since each depends on machinery this branch does not have and none of them is what the change is for: - The statistics registry threading. StatisticsContext here has no registry field and compute never consults providers, so the context is built empty. The memoization is unaffected. - The EnsureRequirements::optimize_with_context override, whose only purpose upstream was to reach that registry. - Phase 0's InterleaveExec normalization, which is context in the upstream diff rather than part of this change. Upstream's ensure_distribution_uses_context_statistics_registry test is dropped for the same reason: it exercises the registry path above. The fetch tests this branch carries from #75 are kept over the upstream versions they conflicted with; they are unrelated to this change. Verified: physical_optimizer::enforce_distribution passes 82/82, clippy and fmt clean. physical_optimizer::sanity_checker has 9 failures, which reproduce on a clean branch-55 checkout and are not from this change.
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
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
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.
Problem
get_repartition_requirement_statuscreates a freshStatisticsContext::new()once per child.StatisticsContext::computerecurses the child's whole subtree and carries a pointer-keyed memoization cache its own docstring calls a "per-call memoization cache" meant to be reused across a walk. Allocating a new context per child throws that cache away every time, so a singleensure_distributionpass recomputes shared subtree statistics O(depth) times over a wide plan.Change
Share one
StatisticsContextacross the wholeensure_distributiontransform_up, passing&StatisticsContextintoget_repartition_requirement_status.The cache is keyed by raw node pointer.
ensure_distributionreturnsTransformed::yesunconditionally, so instead of keying the reset on that flag, the reset fires only when a node's plan pointer actually changed (Arc::as_ptrbefore/after). A changed node may have freed a cached child (which would make a stale pointer key unsafe); an unchanged node cannot, so the cache safely persists across the no-op nodes that dominate a wide plan.Impact
Measured on Atlas's
/stocks/last-trade/vX/snapshotendpoint (EnsureRequirements runs 6x per request, warm):Plan output is unchanged (pure memoization).
Correctness
datafusion --test core_integration physical_optimizer: 541 passed, 0 failed.datafusion-physical-planstatistics tests: 86 passed.Affects
DataFusion physical optimizer only; no plan-shape change, so no behavioral blast radius. Adopted by Atlas via a separate DF rev-bump. Candidate to forward to apache/datafusion (the unused per-call cache is a latent upstream perf miss).
Draft: parent investigation is X-3736; opening early for review while the Atlas rev-bump is prepared.