edmondop opened a new pull request, #25798:
URL: https://github.com/apache/datafusion/pull/25798

   ## Which issue does this PR close?
   
   Related to 
[apache/datafusion#23194](https://github.com/apache/datafusion/issues/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. Scheduling 
belongs to the caller; boundaries contain no stage numbers.
   
   Run an example with:
   
   ```sh
   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 observes pool headroom before priming. Each buffered 
batch still requires a successful reservation; the example reports allocation 
failures as stream errors and does not spill.
   
   ## 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 cover readiness 
across partitions, sequencing dependent boundaries, and memory admission and 
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.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to