apache/datafusion-ballista

Apache DataFusion Ballista Distributed Query Engine

View on GitHub ↗Jump to charts ↓Open shareable report

Summary Information

Updated 35 minutes ago
Added to GitGenius on January 3rd, 2025
Created on May 19th, 2022
Open Issues & Pull Requests: 192 (+0)
Number of forks: 309
Total Stargazers: 2,115 (+0)
Total Subscribers: 45 (+0)

Repository Insights (GitGenius)

Median issue/PR response: 20.2 hours
Mean response time: 218.1 days
90th percentile: 1056.4 days
Tracked items: 360

How this project is maintained

Around half of the issues opened in the past year never receive a reply. 48% of open issues come from outside the core team, a mix of external reports and the maintainers' own roadmap. Work labelled "bug" is answered fastest, typically in about 18 hours, while "enhancement" waits about 5 days. 14% of tracked open issues have had no activity in three months. Only 7% of issues opened in the past year have been closed.

Charts & Analytics

Fetching additional details & charts...

Issue Activity (beta)

Open issues: 148
New in 7 days: 5
Closed in 7 days: 2
Avg open age: 608 days
Stale 30+ days: 122
Stale 90+ days: 76

Recent activity

Opened in 7 days: 5
Closed in 7 days: 2
Comments in 7 days: 3
Events in 7 days: 14

Top labels

  • enhancement (395)
  • bug (182)
  • good first issue (39)
  • help wanted (37)
  • documentation (15)
  • performance (10)
  • TUI (9)
  • development-process (7)

Detailed Description

Apache DataFusion Ballista is a distributed query execution engine built in Rust that extends Apache DataFusion by enabling parallelized execution of workloads across multiple nodes. The project allows existing DataFusion applications to be distributed with minimal code changes, making it accessible for users who want to scale their query processing without major architectural rewrites.

The Ballista architecture consists of scheduler processes and executor processes that can run as native binaries or Docker containers, with deployment options including Docker Compose and Kubernetes. Clients submit jobs to the scheduler, which coordinates task distribution to executors that report back on task status and completion. The system is designed to handle complex SQL queries including CTEs, joins, and subqueries at scale, though the project documentation acknowledges an ongoing gap between DataFusion and Ballista functionality that the community is actively working to close.

Performance benchmarks derived from TPC-H queries demonstrate significant optimization progress. Testing at scale factor 100 with 100 GB of data on a single node with one executor and eight concurrent tasks shows an overall speedup of 2.9x compared to Apache Spark. Individual query performance varies, with some queries showing substantially higher relative speedups than others, indicating that optimization efforts have been particularly effective in certain query patterns.

The codebase is organized into multiple Cargo feature-gated components. The ballista client crate includes a standalone mode feature for in-process scheduler and executor operation. The ballista-core crate provides Arrow IPC optimizations for shuffle performance and optional Spark compatibility mode. The ballista-scheduler component supports optional features including Substrait plan support, Prometheus metrics collection, execution graph visualization, Kubernetes Event Driven Autoscaling integration, REST API endpoints, and stage plan caching control. The ballista-executor crate includes the mimalloc memory allocator for performance optimization alongside Arrow IPC improvements. The ballista-cli component provides a terminal user interface for REST client interactions.

The project shares contributors with apache/datafusion, pingcap/tidb, and nvidia/cudf-spark, indicating cross-pollination with other major distributed data processing systems.