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 #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.

@alamb

alamb commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

@avantgardnerio what do you think about this PR?

@alamb

alamb commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

I think if we are going to put something like this in the Datafusion core repo, I would like to know we had users / consumers of it downstream. For example would this be usable by Ballista (and would it be migrated over)? I want to make sure what goes into the core repo is general purpose and can be used by many projects

Thank you for this PR BTW @edmondop -- it looks well designed

@comphead

comphead commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

I think if we are going to put something like this in the Datafusion core repo, I would like to know we had users / consumers of it downstream.

Yeah, this more like distributed concepts, stages defines narrow transformations, which doesn't require shuffle, the notion that currently doesn't exist in DF core, it would be super useful for datafusion-ballista, datafusion-distributed, however since those 2 downstream projects are independent, so having common component in the core kind of makes sense

@avantgardnerio

avantgardnerio commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Yes, I think Ballista would use it. Our range repartition operators (ORRE/URRE, apache/datafusion-ballista#2169, #2123) take their cuts from a RuntimeStatsExec sketch on the first batch they see. That's only a full sample because a SortExec between them acts as the barrier. Without one, the router sees a single batch. BufferExec (apache/datafusion-ballista#2095) works around this with memory-pool pressure, but that's a stand-in for exactly the prime > ready > inspect > release contract this PR defines.

With StageBoundary we could repartition by range inside a stage without a sort or a pool-pressure heuristic. That's also the gap on the DataFusion side: RepartitionExec routes Partitioning::Range today, but its split points (and the samples from #24766) are fixed at plan time. Computing them from runtime data (#23093) needs a barrier like this one.

I will give the PR a deeper review shortly, but this would allow Datafusion to run parallel windows and all the other use-cases I outlined in #23194

@andygrove

Copy link
Copy Markdown
Member

I haven't looked at PR details yet, but I have pushed for something like this in the past. It would be great for core DF to have this concept and I believe it will make it easier to map to distributed engines such as Ballista

@andygrove

Copy link
Copy Markdown
Member

Following up on my comment above, I went back and found the prior work I was remembering:

  • apache/arrow#7975 made the physical planner replaceable to support distributed execution.
  • apache/arrow#8283 prototyped a DataFusion scheduler that split physical plans into a DAG of stages at partitioning boundaries.
  • datafusion#62 and datafusion#64 proposed extensible execution and a core scheduler that Ballista could extend.
  • The original Ballista stage/shuffle model was developed through #456, #459, #543, #633, #634, #707, #712, #727, #738, and #750.
  • datafusion#3949 moved physical-plan serialization from Ballista into DataFusion.
  • datafusion#11070 proposed a datafusion-distributed crate with common shuffle primitives.
  • datafusion#23282 modeled distributed execution in-process by splitting plans into stages, serializing them, and running partitions as isolated tasks.
  • The Ballista AQE work in #1987, #1988, #1989, and #2092 uses completed stages and runtime statistics to reconsider planning decisions.
  • datafusion-ballista#2434 is a recent concrete example: stage a join’s build side, inspect its exact statistics, and then choose the join strategy.

That history is why I like the direction here. A small core contract for priming a boundary, observing readiness, and releasing buffered output could provide a common seam while leaving scheduling, admission control, materialization, and replanning policy outside core DataFusion.

@Samyak2

Samyak2 commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

This abstraction is very interesting and relevant for us (at e6data).
We have built a similar thing internally, but something "native" in datafusion core would be great to have.

Instead of describing what we're using it for, let me list the capabilities we want to build on top of this:

  1. Being able to schedule one stage at a time
    • Currently, executing a DataFusion plan causes most operators to execute and return a stream.
    • The scheduling is left entirely to Tokio. I understand that this is by design and won't be changed in DataFusion.
    • Having said that, being able to delay execution of an operator will unlock the properties I mention below and also provide more control over how queries are scheduled (for those who want to do that).
  2. Re-planning: opportunity to update the downstream plan (parent operators):
    • This means an abstraction that allows changing what operators sit above the boundary.
    • Ideally these updates should have the same power as the optimizer rules, since those operators have not started yet.
    • So we should be able to add operators, remove operators, change properties of an operator (partitioning, ordering, etc.)
  3. Works well with spill-to-disk
    • When we have boundaries on operators which spill, for example AggregateExec in Final mode, the abstraction should not be holding on to raw RecordBatchs (and then release or spill).
    • It should re-use the operator's spilling capability. For aggregate, this would be completing the stage once all input is done and continue unspilling + aggregating in the next stage.

Looking at these in the context of the specific abstraction introduced in this PR:

  1. Scheduling: StageBoundary makes this more natural. An example usage would be adding a StageBoundary on top of every "pipeline breaker" (like final agg) and only running the stages through this API, not the normal plan.execute(..) API.
  2. Re-planning: StageBoundary also makes this more natural. The ExecutionPlan that implements this can also change its properties at runtime, and then re-plan the tree above it. Again, this is only possible with some external orchestrator API, not the normal plan.execute(..) API.
  3. Spill-to-disk integration: this would be the trickiest part. For this, the operator itself would need to implement StageBoundary. For example, AggregateExec would need to support this new method of execution by implementing StageBoundary itself. This level of control into spilling is not possible by some wrapper operator that sits above. It seems like the intended usage of this API is a wrapper operator, so I think this capability is not satisfied by the proposed API.

Overall: this looks like the right direction. Though, a lot needs to be built on top of this to unlock new capabilities (re-planning, better scheduling).

On the topic of whether this needs to be in datafusion core, or some external crate: things like No. 3 I mentioned above are only possible with deep integration with existing operators. If there's not much community interest in such deep integrations/perf optimizations, then I believe it makes sense to keep this a separate crate and bring it in later if there's enough usage. Also, if we need capabilities like AQE in datafusion core (#23194) (even for single-node execution), then I think some API like this is a pre-requisite.

Personally, I would be very interested in helping out here - discussions, reviews, code or whatever else is needed to push this forward!

@avantgardnerio

Copy link
Copy Markdown
Contributor

Having reviewed this more in depth, I think this is most useful within DataFusion itself.

In Ballista we would make BufferExec implement this, and we'd have the executor stream driver release these when we did unordered range repartitions.

What would be more useful though would be a method on ExecutionPlan, defaulting to None, that existing pipeline breakers like SortExec override:

fn as_boundary(&self) -> Option<&dyn StageBoundary>

A driver could see if something could be treated as a stage boundary, and then optionally prime it as one without the need to wrap or double buffer.

But the real killer use-case (and I'd love to see a demo) would be a conditional release if criteria were met. That would need is_ready to mean "safe to inspect" rather than only EOF.

This would allow DataFusion (irrespective of dfd or Ballista) to do things such as dynamic build-side selection for joins. It could run both legs of a join up until some threshold, and then inspect them to see if either one completed under that threshold. If so it could use that leg as the build side and flip if necessary. I suspect this might help improve DF in real-world benchmarks and make the argument that this belongs in-repo.

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.

7 participants