Repository navigation
Integrate datafusion-distributed with datafusion-python #1612
Description
Activity
There are a couple of PRs upstream trying to adapt the
QueryPlannerto something that can be provided from withindatafusion-python:- refactor: thread SubqueryContext explicitly through physical planning datafusion#22340
- feat: Add
FFI_QueryPlannerto support foreign query planners across shared-library boundaries datafusion#22151
I'm imagining that this is because the idea is that external systems like
datafusion-distributedorballistacan inject their ownQueryPlannerimplementations from within Python (not 100% sure about this, correct me if I'm wrong).But I'm wondering if it might actually be simpler just handle all that details within this project instead.
Thanks for the PR!
The main issue I see with this is that it makes datafusion-distributed a new dependency for datafusion-python. That's going to add bloat to the existing large wheels we're already producing. Also if we want to support ballista in the same way then we we're adding yet another external dependency and trying to ship/support them in this main repo.
The big advantage of this PR is how small / easy it is.
The longer term version I had in mind was that we expose via FFI the physical optimizer (done) and query planner (in progress). Then you have a relatively thin
datafusion-distributedpython package and when you create a session context you simply add in the new query planner or physical optimizer. This would work the same for both ballista, datafusion-distributed, and any other package that comes along and wants to do something similar.What do you think?
Reacted by Andy GroveCompiling the wheels locally with
maturin build --releasegives something like this:Build Wheel size Native .souncompressedorigin/main40.45 MiB 116.41 MiB distributed branch 42.34 MiB 121.57 MiB Increase +1.89 MiB +5.16 MiB So it's like a bit less than a 5% increase in size.
As someone who does not deal with the consequences of big wheels, it seems like a cheap price to pay for the simplicity and the integrated experience for
datafusion-pythonconsumers (they would not need to import an external Python package).But again, it's very easy for me to say this if I'm not the one accountable for the consequences.
🤔 Could there be a way of getting the best of both worlds?
Gathering some thoughts and learnings from a private conversation with @timsaucer. There are two diametrically opposite ways of doing this:
Option 1 — Bundled feature build
datafusion-distributedis added as a Rust dependency ofdatafusion-python, gated behind an off-by-default cargo feature, and compiled into the same wheel. There's no ABI boundary, so the internal wiring (distributed planner, worker server, codecs) is done in plain Rust exactly as the library expects. Users opt in by installing the feature-enabled build; those who don't pay nothing. The cost is thatdatafusion-pythonthen is responsible for maintaining all the wiring. It would also imply treatingdatafusion-distributedas "special", as it's not scalable that all external projects follow this approach.Option 2 — Foreign plugin
datafusion-distributedremains its own separately-built package that plugs intodatafusion-pythonat runtime through datafusion-ffi's stable ABI (PyCapsule).datafusion-pythontakes no dependency on it, keeping the two projects fully decoupled in versioning and governance. The cost is that every interaction must be expressed as an FFI-stable type, and anything the FFI layer can't carry simply can't cross. This is for example the approach Ballista took. The amount of customization optionsdatafusion-distributedoffers to users can make this approach complicated though.The perceived end-result from a user standpoint is very good with Option 1, as distribution would be seamlessly integrated, and it would look like it works "out-of-the-box". The sum of the maintenance burden in all the involved projects (
datafusion,datafusion-distributedanddatafusion-python) is as small as it can get. The complexity of exposing all the customization options ofdatafusion-distributedto users is minimal (e.g., building custom distributed plans injecting network boundaries manually, customizing how custom plans are distributed across workers, etc...).On the other hand, it'd imply treating
datafusion-distributedas special, it would transfer the maintenance burden todatafusion-pythonmaintainers, and it would open the door for a bunch of other external projects to want to do the same.Still, the argument about seamless integration and end-user UX sounds like the most appealing. At the end of the day, I doubt
datafusion-pythonusers care about the ownership model of the different projects involved, most likely they really just want to get the job done with a setup as simple and consistent as possible.So, rather than constructing from the implementation, it might be worth starting from the user experience.
Here's what I think the differences are for the implementation of the options listed above.
Option 1 - Bundled feature build
We have to decide which distributed systems we will support. I suspect we want both
ballistaanddatafusion-distributed. I would prefer to support those that are officially within the DataFusion umbrella. I think DataDog hasdatafusion-distributedstable enough and in production that they are generally willing to donate it to DataFusion.When we have a new major release of DataFusion, we would now need both
ballistaanddatafusion-distributedto make new releases with the upgraded major version before we could do a major release update todatafusion-python. I suspect this would mean that our Python bindings would lag even further behind the upstream repo. We have been trying to keepmainin this repo up to date withmainon DataFusion to try to lower the gap between releases, because this repo also holds up other downstream work from upgrades. My suspicion is that this will add a 2-4 week delay in major releases, based on the way I've seen these release chains work in other projects like geodatafusion, lance, and then rerun (my company's project). Simply the need to have those projects upgrade and go through their own release cycles will slow down releases, not to mention trying to keep up to date with feature upgrades in the datafusion core repo. For me, this is the biggest down side to this option.We now need
datafusion-pythonto add wrapper classes for anything that needs exposure in bothballistaanddatafusion-distributed. Since we already have #1611 we have a very good measure of what the burden is, and a jump start on supporting it. I expectballistato have similar level of effort. We would then need to update our skills to make sure we have coverage on both of those projects in addition to our current upstream coverage checks.The way a user would interact with these would be something along the lines of:
config = SessionConfig().with_distributed(LocalhostWorkerResolver(worker_ports_from_env()))
One issue for
datafusion-distributedspecifically is I'm not 100% sure how to bring a third party resolver to the party. It's been a while since I dug into the code, so maybe it's already simple to do. I think Gabriel said it's not difficult now, so this is possibly a non-issue.Option 2 - Foreign plugin
This would put almost all of the maintenance burden on the individual projects. The work in the
datafusion-pythonrepository is to finish up the upstream query planner work and then to expose these FFI objects. Also I had a branch where we allowed for the default codecs to make it so that well known executors could make it transparently through the FFI boundary, which would be needed bydatafusion-distributed. Basically, we would check if executors are encodable with the default codec. If so, use it instead of making them foreign execs.The down side of this is that now each of these projects would need to maintain the python wrapper code and produce pypi wheels. It increases their burden for release since they are now releasing both rust crates and python wheels, and their developers need to become familiar with all of the PyO3 work, which sometimes requires a lot of changes from release to release.
The way a user would interact with these would be something along the lines of:
ctx = SessionContext().with_query_planner( DataFusionDistributedPlanner(LocalhostWorkerResolver(worker_ports_from_env())) )
There are two issues specific to
datafusion-distributedthat would need to be worked out.- Currently it uses downcasting of executors a lot throughtout the repository. I think we would need to do one of three things: (a) Change these to name based checks, which has the potential to be fickle / not robust (b) categorize each of these checks and add traits to the executor to return some kind of enum that tells if its a shuffle, repartition, etc (c) change the FFI crate to transparently convert these executors, like described above by using the default protobuf codec.
- The current implementation relies on storing data in the config that are not string key/value pairs. I believe this is an anti-pattern anyways and technical debt, but it would have to be addressed right away.
Request
I would like to get other people's opinions besides Gabriel and myself. We're both very interested in making this work, but it's a big decision about how to go about supporting these. Tagging @milenkovicm since you're probably the person most involved on the Ballista side.
thanks for keeping me in the loop @timsaucer ,
IMHO, option 2 makes more sense for few reasons:
- it would not introduce any overhead on datafusion python, which can keep release cadence as it was. currently ballista has release more than 4+ weeks after df has been released. I don't want to put us all on pressure.
- reading the option 1, it discuss something very specific to df-distributed which i don't think it should be concern of df-python, I would suggest to keep separation of concerns
- keeping small codebase locally is not really a problem, as there is not much code. the only issue i see is that each project needs to have release process in place. we have all supporting infrastructure in place already, which could be reused. Py03 is not issue for ballista and i dont think it would be issue for df-distributed, if it is i would be more than happy to help.
- should new framework want to distribute df-python there is a supporting framework in place, and two implementation already using it, making integration framework robust
also, if we move apache/datafusion#22151 to df-python i dont think its would be big effort to get support for
option 2.Currently it uses downcasting of executors a lot throughtout the repository
This is right, downcasting is indeed a really nice offering from
datafusionfrom an API standpoint, and it does help a lot in the codebase, it would be a shame to not count on it.The current implementation relies on storing data in the config that are not string key/value pairs
This is indeed documented tech debt, there was a solution on the table upstream apache/datafusion#18739, but it was not very well received and we rolled it back.
I think now this should be trivial to fix though.
Just to highlight a few disadvantages of Option 2 that are not specific to any particular project:
- It introduces additional constraints on upstream’s internal APIs. For example, refactor: thread SubqueryContext explicitly through physical planning datafusion#22340 imposes the constraint that the session passed to QueryPlanner must be immutable and cannot be cloned. In this particular case, that constraint seems reasonable to me, but it raises the broader question of how well this approach scales if downstream projects continue imposing similar constraints on datafusion.
- From an API and packaging perspective, end users may end up with a more fragmented experience: different APIs inconsistent with their Rust counterparts, managing multiple Python packages (each with overlapping or duplicated wheels), navigating documentation spread across several projects, and manually dealing with version compatibility between those packages.
To me, this second point is the one that should carry the most weight when evaluating Option 1 vs. Option 2, hence the reason we I think it's important to start constructing from end-user experience first.
Regarding
datafusion-distributed, one thing that can be challenging is to bridge all the customization options in the Python world. TBH, I'm not sure which of the two options would make it easier:- Customize worker resolution (done easily here in my attempt to Option 1, I bet this should also be relatively easy with Option 2)
- Customize worker connections (I don't think this one can even be done in any option...)
- Propagating arbitrary headers through worker jumps
- Propagating custom ConfigOptions to workers
- User defined work unit feeds
- Custom worker routing for cache affinity and stuff
- Users injecting network boundary nodes manually in the plan
A lot of these extensibility options are provided as traits. Let's imagine that the trait implementations are carried by SessionConfig.extensions. Does
datafusion-pythonprovide the tools for people to place their custom extensions with arbitrary types there?- From an API and packaging perspective, end users may end up with a more fragmented experience: different APIs inconsistent with their Rust counterparts, managing multiple Python packages (each with overlapping or duplicated wheels), navigating documentation spread across several projects, and manually dealing with version compatibility between those packages.
I'm not sure I understand your concern, with option 2 you'd just need to provide a factory method which would configure
SessionContextto work with df-distributed / ballista, all other behaviour would be the same like df-python,To me, this second point is the one that should carry the most weight when evaluating Option 1 vs. Option 2, hence the reason we I think it's important to start constructing from end-user experience first.
from user perspective it does change much, for example
from datafusion import col, lit from datafusion import DataFrame # we do not need datafusion context # it will be replaced by BallistaSessionContext # from datafusion import SessionContext from ballista import BallistaSessionContext # Change from: # # ctx = SessionContext() # # to: ctx = BallistaSessionContext("df://localhost:50050") # all other functions and functions are from # datafusion module ctx.sql("create external table t stored as parquet location './testdata/test.parquet'") df : DataFrame = ctx.sql("select * from t limit 5") df.show()
ballista support is one line change, all the api is the same, all underlying interaction is with df-python
Reacted by Gabriel and NickThat's quite elegant indeed!
How does that work if you want to add, for example, your custom
ConfigOptions toBallistaSessionContext? or if you want to inject any SessionConfig.extensions?I've not found any examples in https://github.com/apache/datafusion-ballista/tree/main/python
I believe config options should be supported with
FFI_ConfigOptionsnot sure about extensionsDon't expect much from current ballista python, it's a bit of a hack with many rough edges.
+1 for option 2 (foreign plugin)
Option 2 seems best to me. That ballista snippet above looks really clean.
There's a slight benefit to option 1 for discoverability if I could
pip install datafusion[ballista]that makes everything feel a little more unified. However, then datafusion would want to include all downstream options.pip install ballistadoesn't seem that bad if I know I need the distributed support.Managing the python dependencies doesn't seem too bad. I assume the pyproject for ballista and datafusion-distributed would transitively depend on datafusion-python so I don't expect there to be duplicate wheels. If someone needs the distributed option then they wait the couple of weeks after upstream lands, if not then they can get bleeding edge df-python sooner.
👍 Let's give it a go to Option 2 then. AFAIK the next step is to get the following PR in:
I'll review it shortly, cc @timsaucer
- added a commit that references this issue
on Aug 7, 2026 - added a commit that references this issue
on Aug 12, 2026
Is your feature request related to a problem or challenge? Please describe what you are trying to do.
Allow running distributed queries in
datafusion-pythonDescribe the solution you'd like
Ideally something well integrated with
datafusion-pythonthat does not require big changes or using different APIs for executing distributed queries.datafusion-pythonis already a very ergonomic wrapper for usingdatafusion, so something maintaining that philosophy without introducing a lot of API surface would be ideal.I'm interested specifically in using the
datafusion-distributedlibrary from within Python, and I see three mutually exclusive ways of integrating it:datafusion-pythondepend ondatafusion-distributed, hiding some internal plumbing indatafusion-pythonand extending the current API with distributed capabilities.datafusion-distributedanddatafusion-pythonthat ships an external API for using distributed functionality indatafusion-pythondatafusion-distributeddepend ondatafusion-python, providing a set of functions and classes that decoratedatafusion-pythonwith distributed capabilitiesI'm not sure which approach aligns best with this project's philosophy, the naive intuition from someone unfamiliar with this project is that the first option has greater chances of providing a well integrated experience, and it's probably the easiest to implement due to the fact that internal plumbing in the Rust world can be hidden in this project.
I actually tried this here:
And the fact that with only ~1K LOC, examples and tests included, can yield a functional integration, makes me think that it might actually not be a bad idea. But again, I don't know what I don't know, so would very gladly accept feedback and suggestions on something different.