andygrove commented on PR #25798: URL: https://github.com/apache/datafusion/pull/25798#issuecomment-5963872756
Following up on [my comment above](https://github.com/apache/datafusion/pull/25798#issuecomment-5961890786), I went back and found the prior work I was remembering: - [apache/arrow#7975](https://github.com/apache/arrow/pull/7975) made the physical planner replaceable to support distributed execution. - [apache/arrow#8283](https://github.com/apache/arrow/pull/8283) prototyped a DataFusion scheduler that split physical plans into a DAG of stages at partitioning boundaries. - [datafusion#62](https://github.com/apache/datafusion/issues/62) and [datafusion#64](https://github.com/apache/datafusion/issues/64) proposed extensible execution and a core scheduler that Ballista could extend. - The original Ballista stage/shuffle model was developed through [#456](https://github.com/apache/datafusion/issues/456), [#459](https://github.com/apache/datafusion/pull/459), [#543](https://github.com/apache/datafusion/pull/543), [#633](https://github.com/apache/datafusion/pull/633), [#634](https://github.com/apache/datafusion/pull/634), [#707](https://github.com/apache/datafusion/issues/707), [#712](https://github.com/apache/datafusion/pull/712), [#727](https://github.com/apache/datafusion/pull/727), [#738](https://github.com/apache/datafusion/pull/738), and [#750](https://github.com/apache/datafusion/pull/750). - [datafusion#3949](https://github.com/apache/datafusion/issues/3949) moved physical-plan serialization from Ballista into DataFusion. - [datafusion#11070](https://github.com/apache/datafusion/pull/11070) proposed a `datafusion-distributed` crate with common shuffle primitives. - [datafusion#23282](https://github.com/apache/datafusion/pull/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](https://github.com/apache/datafusion-ballista/issues/1987), [#1988](https://github.com/apache/datafusion-ballista/issues/1988), [#1989](https://github.com/apache/datafusion-ballista/issues/1989), and [#2092](https://github.com/apache/datafusion-ballista/issues/2092) uses completed stages and runtime statistics to reconsider planning decisions. - [datafusion-ballista#2434](https://github.com/apache/datafusion-ballista/pull/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. -- 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]
