Skip to content

Serve lazy Datasets over Arrow Flight SQL - #256

Open
alxmrs wants to merge 10 commits into
mainfrom
flight-sql-serve
Open

alxmrs wants to merge 10 commits into
mainfrom
flight-sql-serve

Conversation

@alxmrs

@alxmrs alxmrs commented Sep 26, 2026 •

Copy link
Copy Markdown
Member

Summary

xql.serve starts an Arrow Flight SQL server over the same lazy tables XarrayContext uses. Any Flight SQL client (ADBC's Flight SQL driver in Python/R/Go/Java, the Flight SQL JDBC and ODBC drivers, and SQL tools built on them) can then query a Dataset from another process or machine without copying it.

# server
server = xql.serve({"era5": ds}, host="0.0.0.0", port=8815)
server.wait()

# client, anywhere
import adbc_driver_flightsql.dbapi as flight_sql
con = flight_sql.connect("grpc://server-host:8815")
cur = con.cursor()
cur.execute("SELECT time, AVG(t2m) AS t2m FROM era5 WHERE lat BETWEEN 40 AND 41 GROUP BY time ORDER BY time")
out = xql.to_dataset(cur, template=ds)

Where #255 copies a Dataset into a database, this goes the other way: the data stays where it is, and queries keep partition pruning on dimension predicates and projection pushdown, reading only the chunks and variables they touch, only while they run. Results stream back as Arrow record batches.

Design

  • Own Flight SQL service on arrow-flight's FlightSqlService trait (src/flight.rs), not datafusion-flight-sql-server. That crate unconditionally depends on datafusion-substrait, which needs protoc at build time on every wheel builder, and its datafusion-federation dependency floats to DataFusion 55. arrow-flight ships Flight SQL's generated protobuf code, so the only new dependencies are the gRPC stack (tonic, hyper, prost).
  • Native DataFusion session. A Python SessionContext can't be passed across the FFI boundary into our Rust crate, so the server keeps its own session and registers our PrunableStreamingTable providers directly. Pruning, projection pushdown and exact statistics carry over unchanged.
  • Stateless. Statement tickets and prepared-statement handles are the SQL text, re-planned on use. Implemented: statement and prepared-statement queries, GetCatalogs / GetDbSchemas / GetTables (so adbc_get_objects and SQL tools can browse), and GetSqlInfo.
  • Read-only SQL. DDL, DML and other statements are rejected before planning, so clients can't reach the server's filesystem through CREATE EXTERNAL TABLE or COPY.
  • Safe default bind. 127.0.0.1 unless told otherwise; no auth or TLS, documented.
  • Same register seam. xql.FlightSQLServer() works with xql.register(server, name, ds, table_names=...) via a small adapter, including while it's serving; mixed-dimension Datasets are served as name.group.
  • The server runs on its own thread with a multi-threaded Tokio runtime. Partitions acquire the GIL per Python call, as they do in process, and shutdown(timeout=5.0) releases the GIL while in-flight queries get up to timeout seconds to finish, then closes the remaining connections. (Graceful shutdown alone would wait forever on a client that stopped reading a result partway; caught in review, reproduced, and covered by a test.)

Tested clients

  • Spark 4.0 via the Arrow Flight SQL JDBC driver (19.0.0): full scans round-trip to xarray; aggregations, column subsets, era5.atmosphere, subqueries as dbtable, and Spark SQL over a temp view all match xarray. Spark's filters are fully pushed (*GreaterThanOrEqual(time, ...)) into the SQL it sends, so they reach chunk pruning. The JVM must run in UTC (-Duser.timezone=UTC): otherwise timestamps shift by the JVM zone, twice (+16 h in Pacific). Labeling timestamps UTC server-side isn't a fix, since Spark's JDBC reader rejects timestamp-with-time-zone columns from this driver (unrecognizedSqlTypeError).
  • ClickHouse 26.8 via its arrowFlight table function. It speaks plain Arrow Flight (path descriptors + GetSchema), which arrow-flight's Flight SQL scaffolding rejects, so a thin FlightRouter now serves path descriptors and GetSchema and delegates everything else. A path names a table ('weather', 'era5.surface') or is itself a SELECT/WITH query, which is how ClickHouse, which doesn't push filters into arrowFlight, gets chunk pruning. Time filters need SETTINGS session_timezone = 'UTC' (zone-less times, literals parsed in the server zone).

Known gaps

  • The cftime() and reproject() Python UDFs aren't registered on the server session.
  • No parameter binding for prepared statements, and no transactions (read-only, so nothing to commit).
  • One endpoint per query. Returning one Flight endpoint per group of partitions, so clients can read chunks in parallel (e.g. via adbc_execute_partitions), is the natural follow-up.

Testing

The server reuses the ADBC test harness from #255. In tests/_adbc.py it's one more backend, served: each test starts a server, registers its Datasets there, and queries them through ADBC's Flight SQL driver. It gets the same assertions as every database:

  • ADBC contract: 18 tests (dtypes and NaN, aggregates, name.group for mixed dimensions, timedelta coordinates eager and chunked, time filters, subsecond and nanosecond times, awkward and mixed-case names, non-Latin text coordinates, integer and uint64 extremes, a 10-chunk scan, empty results). Ingest modes and temporary tables skip. Runs in main CI; no service needed.
  • ARCO-ERA5 queries: all 7 (area means, a grid-point series, wind speed, threshold counts, per-level stats, the surface/850 hPa join, a regional round-trip). Run in the in-process job of the adbc databases workflow.
  • Found a bug: a mixed-case table name was folded to lowercase at registration, so "Weather" (quoted) could never be found. Fixed by registering exact names. A mixed-case name now warns that it must be quoted, as on DataFusion over ADBC.
  • ClickHouse as a client in CI: the ClickHouse job now runs the arrowFlight test against a server the test starts, reached through the Docker host gateway. Before, it only ran against a ClickHouse on the same host.
  • tests/test_flight_sql.py keeps what only a server has: registering while serving, rejected DDL/DML/COPY/SET (no file left behind), table discovery, shutdown (including with an unread result), plain Arrow Flight (path as a table or a query, GetSchema), and the Spark and ClickHouse clients. Its round-trip tests moved to the contract.
  • Spark (2 tests, need XARRAY_SQL_TEST_FLIGHT_SQL_JDBC_JAR), passing locally.
  • Full non-integration suite locally: 512 passed. cargo build has no duplicate DataFusion or Arrow crates; clippy (-D warnings) and rustfmt pass.

ADBC's DBAPI warns Cannot disable autocommit on connect, because the server doesn't implement Flight SQL transactions. That's expected for a read-only server and harmless.

🤖 Generated with Claude Code

https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t

@alxmrs
alxmrs force-pushed the flight-sql-serve branch 3 times, most recently from 4085015 to e63b8e3 Compare September 28, 2026 15:42
Base automatically changed from adbc-adapter to main September 28, 2026 17:14
alxmrs and others added 7 commits September 28, 2026 10:15
xql.serve({"era5": ds}) starts a Flight SQL server over the same lazy
tables XarrayContext uses, so any Flight SQL client (ADBC, JDBC, ODBC)
can query a Dataset from another process or machine without copying
it. Queries keep partition pruning and projection pushdown, and results
stream back as Arrow record batches.

- The service is implemented directly on arrow-flight's FlightSqlService
  rather than datafusion-flight-sql-server, which pulls in substrait and
  requires protoc on every wheel builder.
- Stateless: statement tickets and prepared-statement handles are the
  SQL text. Supports catalog/schema/table listings and SqlInfo.
- Read-only: DDL, DML, and other statements are rejected, so clients
  cannot reach the server's filesystem via CREATE EXTERNAL TABLE/COPY.
- Binds 127.0.0.1 by default; no auth or TLS (documented).
- xql.FlightSQLServer works with xql.register like any other engine,
  including registering while serving.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t
Graceful shutdown waits for every open response stream, so a client
that stopped reading a result partway kept shutdown() blocked forever.
shutdown(timeout=5.0) now gives in-flight queries that long, then drops
the server's runtime, closing the remaining connections.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t
ClickHouse's arrowFlight table function (and other plain Flight
clients) names a dataset with a path descriptor and asks for its schema
with GetSchema. arrow-flight's FlightService implementation for a
FlightSqlService decodes every descriptor as a Flight SQL command and
leaves GetSchema unimplemented, so these clients failed with "Not yet
implemented".

A thin router now answers path descriptors and GetSchema and hands
everything else to the Flight SQL service. A path names a table
(["weather"], ["era5.surface"], or one element per part) or is itself a
SELECT/WITH query, which is how a client that cannot push filters down
still gets chunk pruning. Its ticket is an ordinary statement ticket,
so DoGet is unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t
Spark reads xql.serve through the Arrow Flight SQL JDBC driver, and
its filters and column selection reach the server as SQL. The JVM must
run in UTC: over JDBC, timestamps otherwise shift by the JVM zone
(twice, in a non-UTC zone), and labeling them UTC server-side is not
an option because Spark's JDBC reader rejects timestamp-with-time-zone
columns from this driver.

ClickHouse reads through arrowFlight; string literals compared with
zone-less times are parsed in its server zone, so time filters need
session_timezone = 'UTC'.

Spark tests run when XARRAY_SQL_TEST_FLIGHT_SQL_JDBC_JAR is set;
the ClickHouse test when XARRAY_SQL_TEST_CLICKHOUSE_URI is.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t
chDB, which CI now uses for the ClickHouse adapter tests, is built
without the arrowFlight table function, and the test sent queries over
HTTP, so XARRAY_SQL_TEST_CLICKHOUSE_URI=chdb:// crashed it.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t
alxmrs and others added 2 commits September 28, 2026 10:34
DataFusion parses a `&str` table name as SQL, so registering `Weather`
created `weather`, and the quoted `"Weather"` a client writes could never
find it. Mixed-dimension tables already kept their exact names. Register
bare names, and warn on a mixed-case name like the ADBC adapter does for
databases that fold unquoted identifiers.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t
The server joins the tests' backend table as `served`: each test starts
one, registers its Datasets there, and queries them through the ADBC
Flight SQL driver, so the same round-trip assertions (types, missing
values, names, timedelta and subsecond times, mixed dimensions, the
ARCO-ERA5 queries) check the server as they check every database.
Tests of ingest modes and temporary tables skip it. Registration goes
through `Database.register`, which targets the server or the connection.

The Flight SQL tests the contract now covers are removed. The ClickHouse
job also runs ClickHouse's `arrowFlight` against a server the test
starts, reached through the Docker host gateway; before, it only ran
against a ClickHouse on the same host.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t
@alxmrs
alxmrs marked this pull request as ready for review September 29, 2026 00:32
@alxmrs
alxmrs requested a review from Mmoncadaisla September 29, 2026 00:46
A one-part Flight path was parsed as SQL, so `Weather` (as ClickHouse's
arrowFlight sends it) looked for `weather` and missed a table registered
with its exact name. A registered name now matches exactly first, as the
two- and three-part forms already did, falling back to parsing, so
`era5.surface` stays schema-qualified.

An error that stopped the server was dropped, so `wait()` returned as if
it had been shut down. The server thread now ends with that error, and
`shutdown()` and `wait()` raise it.

The mixed-case warning named a frame above the caller: its stack level
was the ADBC adapter's. Each entry point now passes its own.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AMTHTEAoyKzUJLvFmg5G6t

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant