Upgrade Guides#

Ballista 54.0.0#

For library and embedding users#

Upgrade to DataFusion 54#

Ballista 54.0.0 upgrades its DataFusion dependency to 54.0.0. Code built on ballista-core / DataFusion APIs must move with it. See DataFusion’s own Upgrade Guide for 54.0.0 for the DataFusion-side API changes.

job_id is now a newtype#

Scheduler job identifiers are no longer bare Strings. ballista-core now exports a transparent JobId newtype (ballista_core::JobId) so a job id can no longer be confused with another string identifier at the type level. Public signatures that previously carried a String/&str job id now carry a JobId / &JobId — for example DistributedQueryExec::job_id() now returns Option<JobId> instead of Option<String>. Construct one with JobId::new(...) or .into(), and recover the inner string with .as_str() or .into_inner().

// before
let id: Option<String> = query_exec.job_id();

// after
use ballista_core::JobId;
let id: Option<JobId> = query_exec.job_id();
let id_str: Option<String> = id.map(|id| id.into_inner());

For cluster operators#

CORS origins and methods are configured via CLI flags#

The scheduler REST API no longer reads CORS configuration from environment variables. The BALLISTA_CORS_ALLOWED_ORIGINS and BALLISTA_CORS_ALLOWED_METHODS environment variables are removed and replaced by the --cors-allowed-origins and --cors-allowed-methods scheduler CLI flags (each a comma-separated list; * allows any). If you set these environment variables, move the values to the flags — otherwise the scheduler falls back to its defaults (origins http://localhost:8080 and https://nightlies.apache.org; methods GET, PATCH, OPTIONS).

# before (environment variables)
BALLISTA_CORS_ALLOWED_ORIGINS=https://example.com \
BALLISTA_CORS_ALLOWED_METHODS=GET,PATCH \
  ballista-scheduler

# after (scheduler CLI flags)
ballista-scheduler \
  --cors-allowed-origins https://example.com \
  --cors-allowed-methods GET,PATCH

Planner and execution behavior changes#

The static planner now broadcasts small join build sides#

Ballista previously disabled CollectLeft broadcast joins entirely (both datafusion.optimizer.hash_join_single_partition_threshold and hash_join_single_partition_threshold_rows defaulted to 0). In 54.0.0 those defaults are raised to 10 MB (10485760) and 1000000 rows, and the static DefaultDistributedPlanner now converts a SortMergeJoinExec whose smaller side fits under ballista.optimizer.broadcast_join_threshold_bytes (default 10 MB) into a broadcast CollectLeft hash join. This conversion is governed by the new ballista.optimizer.broadcast_sort_merge_join_enabled config key, which defaults to true.

The consequence is that join plans — and therefore performance and shuffle behavior — change on upgrade even with AQE off. Broadcast is applied only to join types that are safe to broadcast (inner and right-side variants); other join types remain repartitioned. To restore the previous behavior, set ballista.optimizer.broadcast_sort_merge_join_enabled = false, ballista.optimizer.broadcast_join_threshold_bytes = 0, datafusion.optimizer.hash_join_single_partition_threshold = 0, and hash_join_single_partition_threshold_rows = 0.

Serialized empty projections#

The plan serialization for a filter that projects to zero columns changed. In 53.x an empty projection (Some([])) and “all columns” (None) both encoded to an empty list and decoded back to None, shifting downstream column indices; the scheduler now rewrites such filters into a serde-safe form before serialization. Together with the DataFusion 54 protobuf bump, this means a 54.0.0 scheduler and a 53.x executor (or vice versa) cannot reliably exchange serialized plans. Upgrade schedulers and executors together — do not run a mixed-version cluster across this upgrade.

AQE join behavior (opt-in)#

If you have enabled ballista.planner.adaptive.enabled = true, adaptive join selection now delays the join decision and can promote a small build side to a broadcast (CollectLeft) hash join at runtime, and empty-join inputs are short-circuited by the propagate-empty rule. This changes the runtime join plans (and their results’ timing/shape) for adaptive clusters; it is gated by ballista.planner.adaptive_join.enabled (default true). Clusters using the default (AQE off) are unaffected.