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]
