Apache DataFusion Ballista: a distributed query engine for DataFusion workloads
Apache DataFusion Ballista Distributed Query Engine. Ballista runs the same SQL and DataFrame workloads across a cluster with minimal code changes and the same results.
At a glance
- What is it?
- Ballista splits DataFusion plans across a scheduler and executors so the same SQL and DataFrame code runs on a cluster. The client swap is small, but the README itself warns about a compatibility gap between DataFusion and Ballista.
- Who is it for?
- Adopt Ballista when you already run DataFusion on one machine and need the same plans spread over several executors, or when you want reusable scheduler and executor building blocks for your own engine. Do not adopt it if you need a project with a documented rollback path, or if your workload fits on one node.
- Can I use it commercially?
- Yes. Apache-2.0 is a permissive licence: you can use, modify and sell software built on it, as long as you keep its copyright and licence notices.
- Is it still maintained?
- Yes. The repository last received commits 5 days ago.
- What is it written in?
- Mainly Rust, according to GitHub's language statistics.
Answers come from the project's GitHub data, last synced on September 25, 2026, and from our analysis. They are not legal advice.
Editorial analysis
The problem Ballista solves for single-node DataFusion users
DataFusion runs a query plan inside one process. That is enough until a scan or a shuffle no longer fits on the machine you have. Ballista takes the same SQL and DataFrame workloads and runs them across a cluster, and the README frames the change as a few lines of code rather than a rewrite. The audience is named explicitly in the README: DataFusion users who have outgrown a single machine, Spark users who want a Rust-native engine with a familiar execution model, and library authors building a specialized engine who want scheduler, executor and plan-serialization building blocks instead of writing distributed execution from scratch. That third group matters. Ballista is published as several crates (ballista, ballista-core, ballista-scheduler, ballista-executor, ballista-api-types, ballista-cli), so it is both an end product and a set of parts. If you only need to run SQL on one box, none of this applies to you.
Scheduler, executors and how a query is split
A cluster is one or more scheduler processes plus one or more executor processes. Clients submit jobs to the scheduler; executors fetch tasks from the scheduler and report task status back. The README points to the architecture guide for detail and describes the interaction in exactly those terms, without publishing a wire protocol in the top-level document.
The execution model is the part Spark users will recognize. Plans are split into stages at shuffle boundaries, one task runs per partition, executors are sized in vcores, and adaptive query execution is part of the model. In docker-compose.yml the executor is started with --vcores 4 and --memory-pool-size 8GB, and the compose file asks for two replicas, so the local topology is one scheduler and two executors. The scheduler is started with --bind-host 0.0.0.0 and --external-host ballista-scheduler, which is the address executors and clients are told to use.
Shuffle performance is a Cargo feature decision, not a runtime one. ballista-core enables arrow-ipc-optimizations by default, described as Arrow IPC optimizations for better shuffle performance. Turning default features off is a build-time choice with a runtime cost.
Installing Ballista and running a first distributed query
The README says the easiest start is one of the standalone or distributed examples under examples/, and then the Getting Started Guide at ballista/client/README.md. There is also a docker-compose.yml at the repository root that brings up a scheduler and two executors using the published images apache/datafusion-ballista-scheduler:latest and apache/datafusion-ballista-executor:latest.
Start the cluster from the repository root:
docker compose upThe scheduler listens on port 50050, which the compose file maps to the host. The scheduler command sets --external-host ballista-scheduler so that executors register under the compose service name rather than localhost.
The client-side change is the point of the project. This is the DataFusion program from the README:
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
let ctx = SessionContext::new();
ctx.register_csv("example", "tests/data/example.csv", CsvReadOptions::new())
.await?;
let df = ctx
.sql("SELECT a, MIN(b) FROM example WHERE a <= b GROUP BY a LIMIT 100")
.await?;
df.show().await?;
Ok(())
}And this is the distributed version, with the context swapped and the rest unchanged:
use ballista::prelude::*;
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
let ctx = SessionContext::standalone().await?;
// register the table and run the query as before
Ok(())
}standalone comes from the ballista client crate's standalone feature, which is on by default and starts the required Ballista infrastructure in the background. That is how you validate a plan locally before pointing the same code at a real scheduler.
Once the cluster is up, opening http://localhost:50050 in a browser redirects to a hosted Web TUI with views for jobs, executors, metrics and scheduler information. The README says the Web TUI can also be run locally, and points to the Ballista CLI documentation for that.
The DataFusion gap the README admits to
The README carries a warning block: there is a gap between DataFusion and Ballista which may bring incompatibilities, and the community is working to close it. That single sentence should shape how you evaluate the project. It means the promise of identical results across single-node and distributed execution is the goal, not a guarantee you can assume for every query. Treat any plan that depends on a recently added DataFusion feature as unverified until you run it on both paths.
The version pairing makes the same point. The workspace pins rust-version 1.94.0 and states in a comment that the project should try to follow the DataFusion version, which is an intention rather than a locking mechanism. The workspace version is 54.0.0 while the most recent release listed is 54.1.0. If you depend on a DataFusion feature that landed after the Ballista release you are on, you are on your own.
There is a second, quieter limitation. Features that matter in production are off by default. prometheus-metrics and graphviz-support are non-default features of ballista-scheduler, and substrait is off unless you enable it. A default build gives you a working cluster and a Web TUI, not metrics export. The README does not document rollback or upgrade procedure for a running cluster, so plan for a restart rather than an in-place migration.
Ballista versus Spark for a Rust-native batch engine
The README deliberately positions Ballista next to Spark SQL, and the comparison is about the execution model rather than raw speed. Ballista keeps the model Spark users know: stages split at shuffle boundaries, one task per partition, executors sized in vcores, adaptive query execution. What changes is the implementation language and the dependency footprint. A Spark job carries the JVM, the Spark runtime and its own catalog layer. Ballista runs as native binaries or Docker images and integrates with DataFusion's Rust crates, so a team already writing Rust does not add a second runtime to the stack.
The trade is maturity and ecosystem breadth. Spark has years of connectors, a documented upgrade path and a large body of operational knowledge. Ballista's README documents a compatibility gap with its own upstream, does not document rollback, and leaves metrics behind a non-default Cargo feature. If your team's skills and existing jobs are Spark, the lighter runtime is not by itself a reason to move. If your code is already Rust and already DataFusion, the migration is a context swap plus a scheduler to operate, and that is a much shorter distance.
Licence and the cost of keeping up
Ballista is Apache-2.0, the same licence as DataFusion, with the standard ASF NOTICE and LICENSE files in the repository root. For most adopters that is the least surprising outcome available: permissive, patent-granting, and compatible with the rest of the DataFusion ecosystem. This is a description of the licence identifier, not legal advice; if you redistribute modified binaries, read LICENSE.txt and NOTICE.txt yourself.
The maintenance picture is concrete. The last push to main was on 2026-08-09, the same day as releases 54.1.0 and 54.0.0, with 53.0.0 tagged later that day. The repository is not archived. Release cadence tracks DataFusion closely, which is good for feature parity and expensive for anyone pinned to an older version: the workspace comment about following the DataFusion version tells you that skipping several releases is not a supported path. Budget for a rebuild whenever you move your DataFusion dependency, and expect the Rust toolchain floor (rust-version 1.94.0) to move with it.
Editorial conclusion
Adopt Ballista when you already run DataFusion on one machine and need the same plans spread over several executors, or when you want reusable scheduler and executor building blocks for your own engine. Do not adopt it if you need a project with a documented rollback path, or if your workload fits on one node. Before committing, verify that the DataFusion version you build against matches the gap the README warns about, and confirm the scheduler port and executor resource flags in docker-compose.yml.
Frequently asked questions
How does DataFusion work?
DataFusion executes SQL and DataFrame query plans inside a single process. Ballista extends it by splitting those plans across a cluster of one or more schedulers and executors.
What is DataFusion used for?
The README describes DataFusion as the engine behind existing applications that register tables and run SQL through a SessionContext. Ballista targets those applications once they outgrow a single machine.
What is Apache DataFusion?
Apache DataFusion is the query engine that Ballista builds on, and Ballista is described as a distributed query execution engine that enhances it for parallel execution across multiple nodes.
Official sources
Add this badge to your README
If you maintain this project, the badge below links readers to this analysis and shows its maintenance status from the daily GitHub snapshot. Paste the markdown into your README; add ?metric=license or ?metric=stars to the image URL for a different field.
[](https://hysenlabs.com/projects/apache-datafusion-ballista)