Skip to content

Integrate datafusion-distributed with datafusion-python #1612

Description

@gabotechs

Is your feature request related to a problem or challenge? Please describe what you are trying to do.

Allow running distributed queries in datafusion-python

Describe the solution you'd like

Ideally something well integrated with datafusion-python that does not require big changes or using different APIs for executing distributed queries.

datafusion-python is already a very ergonomic wrapper for using datafusion, so something maintaining that philosophy without introducing a lot of API surface would be ideal.

I'm interested specifically in using the datafusion-distributed library from within Python, and I see three mutually exclusive ways of integrating it:

  • Make datafusion-python depend on datafusion-distributed, hiding some internal plumbing in datafusion-python and extending the current API with distributed capabilities.
  • Create an external crate that depends on both datafusion-distributed and datafusion-python that ships an external API for using distributed functionality in datafusion-python
  • Make datafusion-distributed depend on datafusion-python, providing a set of functions and classes that decorate datafusion-python with distributed capabilities

I'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.

Activity

  1. gabotechs commented on Jun 26, 2026

    @gabotechs
    Author

    There are a couple of PRs upstream trying to adapt the QueryPlanner to something that can be provided from within datafusion-python:

    I'm imagining that this is because the idea is that external systems like datafusion-distributed or ballista can inject their own QueryPlanner implementations 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.

  2. timsaucer commented on Jun 30, 2026

    @timsaucer
    Member

    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-distributed python 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?

  3. gabotechs commented on Jun 30, 2026

    @gabotechs
    Author

    Compiling the wheels locally with maturin build --release gives something like this:

    Build Wheel size Native .so uncompressed
    origin/main 40.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-python consumers (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?

  4. gabotechs commented on Jul 5, 2026

    @gabotechs
    Author

    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-distributed is added as a Rust dependency of datafusion-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 that datafusion-python then is responsible for maintaining all the wiring. It would also imply treating datafusion-distributed as "special", as it's not scalable that all external projects follow this approach.

    Option 2 — Foreign plugin

    datafusion-distributed remains its own separately-built package that plugs into datafusion-python at runtime through datafusion-ffi's stable ABI (PyCapsule). datafusion-python takes 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 options datafusion-distributed offers to users can make this approach complicated though.

  5. gabotechs commented on Jul 5, 2026

    @gabotechs
    Author

    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-distributed and datafusion-python) is as small as it can get. The complexity of exposing all the customization options of datafusion-distributed to 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-distributed as special, it would transfer the maintenance burden to datafusion-python maintainers, 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-python users 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.

  6. timsaucer commented on Jul 6, 2026

    @timsaucer
    Member

    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 ballista and datafusion-distributed. I would prefer to support those that are officially within the DataFusion umbrella. I think DataDog has datafusion-distributed stable 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 ballista and datafusion-distributed to make new releases with the upgraded major version before we could do a major release update to datafusion-python. I suspect this would mean that our Python bindings would lag even further behind the upstream repo. We have been trying to keep main in this repo up to date with main on 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-python to add wrapper classes for anything that needs exposure in both ballista and datafusion-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 expect ballista to 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-distributed specifically 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-python repository 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 by datafusion-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-distributed that would need to be worked out.

    1. 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.
    2. 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.

  7. milenkovicm commented on Jul 6, 2026

    @milenkovicm
    Contributor

    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 .

  8. gabotechs commented on Jul 6, 2026

    @gabotechs
    Author

    Currently it uses downcasting of executors a lot throughtout the repository

    This is right, downcasting is indeed a really nice offering from datafusion from 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.

  9. gabotechs commented on Jul 6, 2026

    @gabotechs
    Author

    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.

  10. gabotechs commented on Jul 6, 2026

    @gabotechs
    Author

    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:

    A lot of these extensibility options are provided as traits. Let's imagine that the trait implementations are carried by SessionConfig.extensions. Does datafusion-python provide the tools for people to place their custom extensions with arbitrary types there?

  11. milenkovicm commented on Jul 6, 2026

    @milenkovicm
    Contributor
    • 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 SessionContext to 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

  12. gabotechs commented on Jul 6, 2026

    @gabotechs
    Author

    That's quite elegant indeed!

    How does that work if you want to add, for example, your custom ConfigOptions to BallistaSessionContext? 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

  13. milenkovicm commented on Jul 6, 2026

    @milenkovicm
    Contributor

    I believe config options should be supported with FFI_ConfigOptions not sure about extensions

    Don't expect much from current ballista python, it's a bit of a hack with many rough edges.

  14. andygrove commented on Jul 6, 2026

    @andygrove
    Member

    +1 for option 2 (foreign plugin)

  15. ntjohnson1 commented on Jul 7, 2026

    @ntjohnson1
    Contributor

    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 ballista doesn'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.

  16. gabotechs commented on Jul 8, 2026

    @gabotechs
    Author

    👍 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

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or request

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions