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 |
|
URL |
|
Advanced options |
|
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-runCREATE EXTERNAL TABLEon each connection. See Extending Ballista components.
Limitations#
These are known gaps, tracked in #2298:
GetFlightInfoblocks until the query finishes.PollFlightInfois 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
DoPutPreparedStatementQueryparameter binding is not, socur.execute(sql, parameters=...)will fail.No query cancellation. A client only holds a query’s
FlightInfoonceGetFlightInfohas 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, andCOPYare rejected with a clear error rather than silently executing on the scheduler, including when wrapped inEXPLAIN ANALYZE. So areCREATE TABLE AS SELECTand SQLEXECUTE, which would otherwise run their query on the scheduler rather than the cluster; use a Flight SQL prepared statement instead ofEXECUTE. Other DDL is supported, includingCREATE EXTERNAL TABLEandCREATE VIEW.No transactions or savepoints.
No Substrait.
CommandStatementSubstraitPlanis not implemented, even when the scheduler’ssubstraitfeature is enabled.Primary/foreign key metadata is not implemented, so tools that browse relationships will show none.