# arrow-ballista
**Repository Path**: mirrors_apache/arrow-ballista
## Basic Information
- **Project Name**: arrow-ballista
- **Description**: Apache DataFusion Ballista Distributed Query Engine
- **Primary Language**: Unknown
- **License**: Apache-2.0
- **Default Branch**: main
- **Homepage**: None
- **GVP Project**: No
## Statistics
- **Stars**: 0
- **Forks**: 0
- **Created**: 2022-06-16
- **Last Updated**: 2026-07-18
## Categories & Tags
**Categories**: Uncategorized
**Tags**: None
## README
# Ballista: Making DataFusion Applications Distributed
[![Apache licensed][license-badge]][license-url]
[license-badge]: https://img.shields.io/badge/license-Apache%20v2-blue.svg
[license-url]: https://github.com/apache/datafusion-comet/blob/main/LICENSE.txt
Ballista is a distributed query execution engine that enhances [Apache DataFusion](https://github.com/apache/datafusion)
by enabling the parallelized execution of workloads across multiple nodes in a distributed environment.
Existing DataFusion application:
```rust
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
let ctx = SessionContext::new();
// register the table
ctx.register_csv("example", "tests/data/example.csv", CsvReadOptions::new())
.await?;
// create a plan to run a SQL query
let df = ctx
.sql("SELECT a, MIN(b) FROM example WHERE a <= b GROUP BY a LIMIT 100")
.await?;
// execute and print results
df.show().await?;
Ok(())
}
```
can be distributed with few lines of code changed:
> [!IMPORTANT]
> There is a gap between DataFusion and Ballista, which may bring incompatibilities. The community is actively working
> to close the gap
```rust
use ballista::prelude::*;
use datafusion::prelude::*;
#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
// create SessionContext with ballista support
// standalone context will start all required
// ballista infrastructure in the background as well
let ctx = SessionContext::standalone().await?;
// everything else remains the same
// register the table
ctx.register_csv("example", "tests/data/example.csv", CsvReadOptions::new())
.await?;
// create a plan to run a SQL query
let df = ctx
.sql("SELECT a, MIN(b) FROM example WHERE a <= b GROUP BY a LIMIT 100")
.await?;
// execute and print results
df.show().await?;
Ok(())
}
```
For documentation or more examples, please refer to the [Ballista User Guide][user-guide].
## Who is Ballista for
Ballista serves several distinct audiences:
- **DataFusion users going multi-node** — you already use [Apache DataFusion](https://github.com/apache/datafusion) on a single machine and have outgrown it. Ballista runs the same SQL and DataFrame workloads across a cluster with minimal code changes and the same results.
- **Spark users wanting the same execution model** — you run Spark SQL or batch jobs and want a lighter, Rust-native alternative without relearning a new paradigm. Ballista keeps the familiar model: plans split into stages at shuffle boundaries, one task per partition, executors with vcores, and adaptive query execution (AQE).
- **Library users building a specialized engine** — you are building a bespoke distributed query engine and want reusable scheduler, executor, and plan-serialization building blocks with extension points, instead of writing distributed execution from scratch.
These audiences are documented in more detail, along with the guarantees each relies on, in the [User Personas](docs/source/contributors-guide/user-personas.md) guide.
## Architecture
A Ballista cluster consists of one or more scheduler processes and one or more executor processes. These processes
can be run as native binaries and are also available as Docker Images, which can be easily deployed with
[Docker Compose](https://datafusion.apache.org/ballista/user-guide/deployment/docker-compose.html) or
[Kubernetes](https://datafusion.apache.org/ballista/user-guide/deployment/kubernetes.html).
The following diagram shows the interaction between clients and the scheduler for submitting jobs, and the interaction
between the executor(s) and the scheduler for fetching tasks and reporting task status.

See the [architecture guide](docs/source/contributors-guide/architecture.md) for more details.
## Getting Started
The easiest way to get started is to run one of the standalone or distributed [examples](./examples/README.md). After
that, refer to the [Getting Started Guide](ballista/client/README.md).
## Cargo Features
Ballista uses Cargo features to enable optional functionality. Below are the available features for each crate.
### ballista (client)
| Feature | Default | Description |
| ------------ | ------- | -------------------------------------------------------------- |
| `standalone` | Yes | Enables standalone mode with in-process scheduler and executor |
### ballista-core
| Feature | Default | Description |
| ------------------------- | ------- | ---------------------------------------------------------------------- |
| `arrow-ipc-optimizations` | Yes | Enables Arrow IPC optimizations for better shuffle performance |
| `spark-compat` | No | Enables Spark compatibility mode via datafusion-spark |
| `build-binary` | No | Required for building binary executables (AWS S3 support, CLI parsing) |
| `force_hash_collisions` | No | Testing-only: forces all values to hash to same value |
### ballista-scheduler
| Feature | Default | Description |
| -------------------------- | ------- | ------------------------------------------------ |
| `build-binary` | Yes | Builds the scheduler binary with CLI and logging |
| `substrait` | No | Enables Substrait plan support |
| `prometheus-metrics` | No | Enables Prometheus metrics collection |
| `graphviz-support` | No | Enables execution graph visualization |
| `spark-compat` | No | Enables Spark compatibility mode |
| `keda-scaler` | No | Kubernetes Event Driven Autoscaling integration |
| `rest-api` | No | Enables REST API endpoints |
| `disable-stage-plan-cache` | No | Disables caching of stage execution plans |
### ballista-executor
| Feature | Default | Description |
| ------------------------- | ------- | ----------------------------------------------------- |
| `arrow-ipc-optimizations` | Yes | Enables Arrow IPC optimizations |
| `build-binary` | Yes | Builds the executor binary with CLI and logging |
| `mimalloc` | Yes | Uses mimalloc memory allocator for better performance |
| `spark-compat` | No | Enables Spark compatibility mode |
### ballista-cli
| Feature | Default | Description |
| ------- | ------- | -------------------------------------------------- |
| `tui` | Yes | Enables a REST client with Terminal User Interface |

### Usage Examples
```bash
# Build with standalone support (default)
cargo build -p ballista
# Build with Substrait support
cargo build -p ballista-scheduler --features substrait
# Build with Spark compatibility
cargo build -p ballista-executor --features spark-compat
```
## Project Status
Ballista supports a wide range of SQL, including CTEs, Joins, and subqueries and can execute complex queries at scale,
but still there is a gap between DataFusion and Ballista which we want to bridge in near future.
Refer to the [DataFusion SQL Reference](https://datafusion.apache.org/user-guide/sql/index.html) for more
information on supported SQL.
## Who uses Ballista
The following organizations use Ballista. To add yours, open a pull request.
| Organization | |
| ------------------------------------------------------------------------------------------------------------------------------ | ------------------------------------------------------------- |
|
| [Spice AI](https://spice.ai/blog/apache-ballista-at-spice-ai) |
|
| [Coralogix](https://coralogix.com/) |
## Contribution Guide
Please see the [Contribution Guide](CONTRIBUTING.md) for information about contributing to Ballista.
[user-guide]: https://datafusion.apache.org/ballista/