Skip to content

[X-3748] ensure_distribution: reuse one StatisticsContext across the pass - #78

Merged
zhuqi-lucas merged 3 commits into
branch-55from
qizhu/x-3748
Sep 18, 2026
Merged

zhuqi-lucas merged 3 commits into
branch-55from
qizhu/x-3748

Conversation

@zhuqi-lucas

Copy link
Copy Markdown

Problem

get_repartition_requirement_status creates a fresh StatisticsContext::new() once per child. StatisticsContext::compute recurses 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 single ensure_distribution pass recomputes shared subtree statistics O(depth) times over a wide plan.

Change

Share one StatisticsContext across the whole ensure_distribution transform_up, passing &StatisticsContext into get_repartition_requirement_status.

The cache is keyed by raw node pointer. ensure_distribution returns Transformed::yes unconditionally, so instead of keying the reset on that flag, the reset fires only when a node's plan pointer actually changed (Arc::as_ptr before/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/snapshot endpoint (EnsureRequirements runs 6x per request, warm):

  • EnsureRequirements: 130.3ms -> 117.3ms
  • physical planning: 156.5ms -> 142.4ms
  • ~13ms per request

Plan output is unchanged (pure memoization).

Correctness

  • datafusion --test core_integration physical_optimizer: 541 passed, 0 failed.
  • datafusion-physical-plan statistics tests: 86 passed.
  • fmt + clippy clean.

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.

…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
zhuqi-lucas marked this pull request as ready for review September 17, 2026 03:43
Copilot AI lite review requested due to automatic review settings September 17, 2026 03:43

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟢 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_stats with caller-provided statistics context.
  • Updated EnsureRequirements to 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.

@zhuqi-lucas
zhuqi-lucas merged commit 1e82195 into branch-55 Sep 18, 2026
69 checks passed
@zhuqi-lucas
zhuqi-lucas deleted the qizhu/x-3748 branch September 18, 2026 08:01
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.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants