Skip to content

fix: preserve fetch across distribution reoptimization - #75

Merged
xudong963 merged 2 commits into
branch-55from
xudong963/preserve-topk-fetch
Aug 27, 2026
Merged

xudong963 merged 2 commits into
branch-55from
xudong963/preserve-topk-fetch

Conversation

@xudong963

Copy link
Copy Markdown

Which issue does this PR close?

  • N/A (fork-specific regression found during the Atlas DataFusion 55 upgrade).

Rationale for this change

Repeated EnsureRequirements passes can remove a distribution operator that carries fetch without restoring it. A query such as ORDER BY ... LIMIT 10 can therefore lose its limit, return all rows, and turn a bounded Top-K into an unbounded allocation.

What changes are included in this PR?

  • Preserve the minimum fetch while distribution-changing operators are removed and rebuilt.
  • Move fetch from a replaced sort-preserving merge to the replacement sort when ordering must be restored.
  • Add regressions for both optimizer paths.

Are these changes tested?

Yes. The full physical_optimizer::enforce_distribution integration-test module passes (81 tests), and datafusion-physical-optimizer passes clippy with warnings denied.

Are there any user-facing changes?

Queries retain their requested LIMIT across physical optimizer passes. There are no public API changes.

@xudong963
xudong963 requested a review from zhuqi-lucas August 27, 2026 09:44
@xudong963
xudong963 merged commit 1556abc into branch-55 Aug 27, 2026
69 checks passed
@xudong963
xudong963 deleted the xudong963/preserve-topk-fetch branch August 27, 2026 10:37
zhuqi-lucas added a commit that referenced this pull request Sep 17, 2026
…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 added a commit that referenced this pull request Sep 18, 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.
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.
MassivePizza added a commit that referenced this pull request Sep 30, 2026
Brings in the upstream 55.1.0 backports. Content matches rebasing the fork
patches onto upstream/branch-55, with these fork commits superseded upstream:

- #75 fix: preserve fetch across distribution reoptimization
  -> apache#24809 (backport apache#25821)
- #79 fix(ci): pull the MinIO test image from quay.io
  -> apache#25092, apache#25216 (backport apache#25485), later replaced by
     the switch to rustfs in apache#25706 (backport apache#25759)

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
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