Samyak2 commented on PR #25798:
URL: https://github.com/apache/datafusion/pull/25798#issuecomment-6014320580
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
`RecordBatch`s (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
(https://github.com/apache/datafusion/issues/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!
--
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]