Arrow Flight SQL#

Ballista can serve Arrow Flight SQL directly from the scheduler, so any Flight SQL client can run distributed queries without linking against Ballista: Python’s ADBC driver, the Arrow Flight SQL JDBC driver, and BI tools that speak either.

Clients send SQL text. The scheduler plans it, distributes execution across the cluster, and returns one Flight endpoint per output partition. Results are streamed back through the scheduler, so clients never need to reach executors and the frontend works unchanged behind NAT, Docker, Kubernetes, and load balancers.

Enabling it#

Flight SQL is a non-default compile-time feature and is also off at runtime, so a scheduler never starts serving SQL by accident.

cargo build --release -p ballista-scheduler --features flight-sql
RUST_LOG=info ./target/release/ballista-scheduler --flight-sql

Start one or more executors as usual:

RUST_LOG=info ./target/release/ballista-executor -c 4 -p 50051

Flight SQL is served on the scheduler’s existing gRPC port (50050 by default) — there is no second port to expose.

Warning

Ballista ships no authentication. With --flight-sql and no authenticator, any client that can reach the scheduler port can run queries, and all unauthenticated clients share a single session (and therefore a single catalog). The scheduler logs a warning at startup when this is the case.

There is a second reason not to expose the port: Flight tickets carry the executor address the scheduler should fetch a partition from, and the scheduler does not verify that the address belongs to the cluster. A client that forges a ticket can make the scheduler open a gRPC connection to an arbitrary host and relay the response. This is a property of the existing Flight result proxy rather than of Flight SQL, but enabling Flight SQL is what makes the port something you would otherwise consider exposing.

Do not expose the port outside a trusted network. See Authentication below.

Querying from Python with ADBC#

Install the driver:

pip install adbc-driver-flightsql pyarrow

Connect to the scheduler and run a query. There is no Ballista Python package involved here — this is the plain ADBC driver talking to the cluster:

import adbc_driver_flightsql.dbapi as flight_sql

with flight_sql.connect("grpc://localhost:50050") as conn:
    with conn.cursor() as cur:
        # Register a table. DDL runs on the scheduler and lands in the
        # session's catalog.
        cur.execute(
            """
            CREATE EXTERNAL TABLE trips
            STORED AS PARQUET
            LOCATION '/data/yellow_tripdata_2022-01.parquet'
            """
        )

        # This one is planned by the scheduler and executed across the cluster.
        cur.execute(
            """
            SELECT payment_type, COUNT(*) AS trips, AVG(total_amount) AS avg_fare
            FROM trips
            GROUP BY payment_type
            ORDER BY trips DESC
            """
        )
        table = cur.fetch_arrow_table()

print(table)

fetch_arrow_table() reads every partition the scheduler advertised, so the result arrives as Arrow the whole way from the executors — no row-by-row conversion.

Introspection#

The driver’s metadata calls are answered from the scheduler’s real DataFusion catalog, so what you see is what you can query:

with flight_sql.connect("grpc://localhost:50050") as conn:
    print(conn.adbc_get_table_types())
    print(conn.adbc_get_objects(depth="tables").read_all())

    with conn.cursor() as cur:
        cur.execute("SELECT 1")
        print(cur.description)

Streaming large results#

fetch_arrow_table() materializes everything. For results that do not fit in memory, take the record batch reader instead and consume it incrementally:

with conn.cursor() as cur:
    cur.execute("SELECT * FROM trips")
    reader = cur.fetch_record_batch()
    for batch in reader:
        ...  # one Arrow RecordBatch at a time

Connecting with JDBC#

Download the Arrow Flight SQL JDBC driver from Maven Central and point your tool at the scheduler:

Setting

Value

Driver class

org.apache.arrow.driver.jdbc.ArrowFlightJdbcDriver

URL

jdbc:arrow-flight-sql://localhost:50050

Advanced options

useEncryption=false

useEncryption=false is required because Ballista serves Flight SQL over plain gRPC; put a TLS-terminating proxy in front of it if you need encryption.

Authentication#

The frontend has no built-in credentials. Ballista’s default authenticator accepts every handshake, which is why the startup warning exists.

To require authentication, implement ballista_flight_sql::Authenticator and pass it to the scheduler through SchedulerConfig. Because the trait sees the request metadata, HTTP Basic credentials (the Flight convention) and bearer tokens both work:

use std::sync::Arc;
use ballista_flight_sql::{Authenticator, Identity};
use ballista_scheduler::config::SchedulerConfig;
use tonic::{Status, metadata::MetadataMap};

struct MyAuth;

#[async_trait::async_trait]
impl Authenticator for MyAuth {
    async fn authenticate(&self, headers: &MetadataMap) -> Result<Identity, Status> {
        // Validate `headers["authorization"]` however your deployment requires.
        Ok(Identity::user("alice"))
    }
}

let config = SchedulerConfig::default()
    .with_flight_sql(true)
    .with_flight_sql_authenticator(Arc::new(MyAuth));

With an authenticator installed, clients must complete the Flight handshake and send the returned bearer token on every Flight SQL request; each handshake gets its own Ballista session, so catalogs, prepared statements, and result tickets are no longer shared.

An authenticator gives each client a session of its own. It is not a security boundary for the cluster: the same port also serves the scheduler’s gRPC API, the REST API (with the rest-api feature), and Ballista’s own partition-fetch tickets, none of which check a token. Keep the port on a trusted network even with an authenticator installed.

Alternatively, terminate authentication in a Tonic interceptor or a proxy in front of the scheduler and leave the frontend as-is.

Sessions and the catalog#

  • A client that completes a handshake gets a private session. Tables it creates with DDL are visible only to that connection.

  • Clients that do not authenticate share one anonymous session.

  • Sessions, prepared statements, and unredeemed result tickets expire after 30 minutes idle.

  • Embedders can pre-populate the catalog for every session through the scheduler’s SessionBuilder, so users do not have to re-run CREATE EXTERNAL TABLE on each connection. See Extending Ballista components.

Limitations#

These are known gaps, tracked in #2298:

  • GetFlightInfo blocks until the query finishes. PollFlightInfo is not implemented, so a long query can hit a client-side deadline. Raise your client’s timeout for TPC-H-scale queries.

  • No bound parameters. Prepared statements are supported, but DoPutPreparedStatementQuery parameter binding is not, so cur.execute(sql, parameters=...) will fail.

  • No query cancellation. A client only holds a query’s FlightInfo once GetFlightInfo has returned, by which point the query has finished. A client that disconnects early leaves its query running to completion.

  • No write path. INSERT, UPDATE, DELETE, and COPY are rejected with a clear error rather than silently executing on the scheduler, including when wrapped in EXPLAIN ANALYZE. So are CREATE TABLE AS SELECT and SQL EXECUTE, which would otherwise run their query on the scheduler rather than the cluster; use a Flight SQL prepared statement instead of EXECUTE. Other DDL is supported, including CREATE EXTERNAL TABLE and CREATE VIEW.

  • No transactions or savepoints.

  • No Substrait. CommandStatementSubstraitPlan is not implemented, even when the scheduler’s substrait feature is enabled.

  • Primary/foreign key metadata is not implemented, so tools that browse relationships will show none.