Skip to content

feat(physical-plan): add staged execution boundary contract - #25798

Open
edmondop wants to merge 5 commits into
apache:mainfrom
edmondop:stage-boundary-api
Open

edmondop wants to merge 5 commits into
apache:mainfrom
edmondop:stage-boundary-api

Conversation

@edmondop

@edmondop edmondop commented Sep 27, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Related to apache/datafusion#23194.

Execution model

DataFusion normally drives a query by polling the output streams returned by ExecutionPlan::execute. Operators pull batches from their inputs; eager operators can also poll inputs in background tasks. This works well for pipelined queries and parallel partitions without requiring a central stage scheduler.

An external system needs additional control when it must finish an input subtree before letting downstream work proceed—for example to inspect a materialized result, admit another memory-intensive stage, or replan the remainder. Although operators such as SortExec already buffer internally, ExecutionPlan has no common completion-and-release interface for coordinating those decisions.

Rationale for this change

StageBoundary provides per-partition priming and readiness, followed by explicit release of buffered output. The caller chooses execution order from the plan dependencies and can inspect completed inputs before releasing them.

What changes are included in this PR?

  • Adds the public trait and a short driver example to datafusion-physical-plan.
  • Adds an InMemoryStageBoundaryExec demonstration in datafusion-examples. It collects batches into a vector per partition and accounts for them with MemoryReservation.
  • Adds three runnable examples: deterministic pause/resume with drain timing, memory-aware admission, and two dependent boundaries.

Run an example with:

cargo run -p datafusion-examples --example execution_monitoring -- stage_pause
cargo run -p datafusion-examples --example execution_monitoring -- stage_admission
cargo run -p datafusion-examples --example execution_monitoring -- stage_dependencies

The admission example waits for a shared memory budget before starting another stage, and returns the budget after the buffered output is consumed.

What is the testing strategy for this PR?

Tests verify that a boundary can finish its input without emitting output, then resume with unchanged partitioned results. They also check that queued admission leaves inputs unstarted until budget is returned, all-partition readiness, dependent boundaries, cancellation of blocked drains, release of reservations when output is dropped, and panic/error propagation. The runnable examples use the same implementation as the tests, and a compiling rustdoc example checks the public API.

Are there any user-facing changes?

Adds a public StageBoundary trait and example commands. Existing ExecutionPlan implementations and SQL behavior are unaffected.

@github-actions github-actions Bot added the physical-plan Changes to the physical-plan crate label Sep 27, 2026
@codecov-commenter

codecov-commenter commented Sep 27, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 82.51%. Comparing base (871058c) to head (351823f).

Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25798      +/-   ##
==========================================
- Coverage   82.51%   82.51%   -0.01%     
==========================================
  Files        1141     1141              
  Lines      439801   439801              
  Branches   439801   439801              
==========================================
- Hits       362916   362908       -8     
- Misses      54943    54947       +4     
- Partials    21942    21946       +4     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@edmondop
edmondop marked this pull request as ready for review September 27, 2026 19:22
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

physical-plan Changes to the physical-plan crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants