Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
102 changes: 102 additions & 0 deletions .github/workflows/ci-wasm.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
name: ci-wasm

on:
push:
branches: [ main ]
pull_request:
branches: [ main ]
workflow_dispatch:
# publish.yml calls this to add the Pyodide wheel to each release.
workflow_call:

concurrency:
group: ${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: true

env:
# The Pyodide release to build for. Its Python version, Emscripten version,
# and Rust toolchain all come from this release's cross-build environment.
PYODIDE_VERSION: "314.0.7"

jobs:
build:
name: Pyodide wheel
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v6

# pyodide-build requires the host Python to match Pyodide's Python.
- uses: actions/setup-python@v5
with:
python-version: "3.14"

- name: Install pyodide-build and the cross-build environment
run: |
pip install pyodide-build
pyodide xbuildenv install "$PYODIDE_VERSION"
echo "RUST_TOOLCHAIN=$(pyodide config get rust_toolchain)" >> "$GITHUB_ENV"

- name: Install Rust
uses: dtolnay/rust-toolchain@master
with:
toolchain: ${{ env.RUST_TOOLCHAIN }}
targets: wasm32-unknown-emscripten

- name: Install Emscripten
run: pyodide xbuildenv install-emscripten

- name: Build
env:
# The Flight SQL server needs sockets and threads; the wasm
# profile optimizes for download size (see Cargo.toml).
MATURIN_PEP517_ARGS: --no-default-features --profile wasm
run: pyodide build -o dist

- uses: actions/setup-node@v4
with:
node-version: "24"

- name: Query a Dataset with DuckDB in Pyodide
run: |
npm install --no-save "pyodide@$PYODIDE_VERSION"
cat > smoke.mjs <<'EOF'
import { loadPyodide } from "pyodide";
import { readdirSync, readFileSync } from "node:fs";
const wheel = readdirSync("dist").find((f) => f.endsWith(".whl"));
const py = await loadPyodide();
await py.loadPackage("micropip");
py.FS.writeFile(`/tmp/${wheel}`, readFileSync(`dist/${wheel}`));
await py.runPythonAsync(`
import micropip
await micropip.install(["emfs:/tmp/${wheel}", "duckdb"])

import importlib.util

import duckdb
import numpy as np
import pandas as pd
import xarray as xr
import xarray_sql as xql

ds = xr.Dataset(
{"air": (("time", "lat"), np.arange(12.0).reshape(4, 3))},
coords={
"time": pd.date_range("2020-01-01", periods=4),
"lat": [10.0, 20.0, 30.0],
},
)
con = duckdb.connect()
xql.register(con, "air", ds, chunks={"time": 2})
rel = con.sql("SELECT time, lat, air FROM air ORDER BY time, lat")
xr.testing.assert_identical(xql.to_dataset(rel, template=ds), ds)
assert not hasattr(xql, "Context")
assert importlib.util.find_spec("dask") is None
`);
console.log(`${wheel} round-trips through DuckDB in Pyodide ${py.version}`);
EOF
node smoke.mjs

- uses: actions/upload-artifact@v6
with:
name: pyodide-dist
path: dist/*.whl
18 changes: 17 additions & 1 deletion .github/workflows/publish.yml
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,12 @@ jobs:
name: dist-manylinux-aarch64
path: dist/*

# Runs the Pyodide smoke test too, so a release never ships an untested
# wasm wheel. PyPI accepts these wheels under PEP 783, so Pyodide users
# can `await micropip.install("xarray-sql")`.
build-artifacts-pyodide:
uses: ./.github/workflows/ci-wasm.yml

build-sdist:
name: Source distribution
runs-on: ubuntu-latest
Expand Down Expand Up @@ -193,7 +199,7 @@ jobs:
uv pip install --extra-index-url https://test.pypi.org/simple --upgrade de
python -c "import xarray_sql; print(xarray_sql.__version__)"
upload-to-pypi:
needs: verify-built-dist
needs: [verify-built-dist, build-artifacts-pyodide]
runs-on: ubuntu-latest
# Only publish for real releases. `workflow_dispatch` runs build and
# verify the wheels (uploaded as artifacts) without pushing to PyPI, so
Expand All @@ -211,3 +217,13 @@ jobs:

- name: Publish package to PyPI
run: uv publish --token ${{ secrets.PYPI_TOKEN }}

# Uploaded separately and after the native wheels, so a rejected wasm
# wheel cannot leave a partial release.
- uses: actions/download-artifact@v6
with:
name: pyodide-dist
path: pyodide-dist

- name: Publish Pyodide wheel to PyPI
run: uv publish --token ${{ secrets.PYPI_TOKEN }} pyodide-dist/*
25 changes: 25 additions & 0 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,31 @@ already on `PATH`) and start at step 2.
You can also run the hooks manually with: `uvx pre-commit run --all-files`
7. Build and serve docs locally: `uvx zensical serve`

### Building for WebAssembly (Pyodide)

The `ci-wasm` workflow builds a Pyodide wheel. The target Pyodide release
fixes the Python version, the exact Emscripten version, and the Rust
toolchain. The Flight SQL server needs sockets and threads, so wasm builds
turn off the `flight` cargo feature.

To build the wheel locally, set `PYODIDE_VERSION` to the value in
`.github/workflows/ci-wasm.yml`. pyodide-build must run on the same Python
version as that Pyodide release (for example, Python 3.14 for Pyodide 314),
which `uv tool install --python` sets up. The `pyodide` command comes from
`pyodide-cli`, so install it with `pyodide-build`:

```shell
uv tool install --python 3.14 pyodide-cli --with pyodide-build
pyodide xbuildenv install "$PYODIDE_VERSION"
rustup toolchain install "$(pyodide config get rust_toolchain)" \
--target wasm32-unknown-emscripten
pyodide xbuildenv install-emscripten
MATURIN_PEP517_ARGS="--no-default-features --profile wasm" pyodide build -o dist
```

On NixOS, run these commands in `nix run .#wasm-shell`. This shell is an
FHS environment, which lets the prebuilt Emscripten and rustup binaries run.


## Before submitting a pull request...

Expand Down
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

27 changes: 22 additions & 5 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -19,22 +19,39 @@ exclude = [

[dependencies]
arrow = { version = "58", features = ["pyarrow"] }
arrow-flight = { version = "58", features = ["flight-sql"] }
arrow-flight = { version = "58", features = ["flight-sql"], optional = true }
async-stream = "0.3"
async-trait = "0.1"
datafusion = { version = "54.0.0" }
datafusion = { version = "54.0.0", default-features = false }
datafusion-ffi = { version = "54.0.0" }
futures = { version = "0.3" }
half = "2.7"
prost = "0.14"
prost = { version = "0.14", optional = true }
# `abi3-py310` builds against CPython's stable ABI, so a single wheel per
# platform works on all CPython >= 3.10 (matching `requires-python`). Maturin
# enables `pyo3/extension-module` through pyproject.toml for wheel builds; it
# must stay disabled for ordinary Cargo test binaries so they link libpython.
pyo3 = { version = "0.28.0", features = ["abi3-py310"] }
tokio = { version = "1.46.1", features = ["macros", "net", "rt", "rt-multi-thread", "sync", "time"] }
tonic = { version = "0.14", features = ["transport"] }
tokio = { version = "1.46.1", features = ["macros", "rt", "sync", "time"] }
tonic = { version = "0.14", features = ["transport"], optional = true }

[features]
default = ["flight"]
# The Arrow Flight SQL server needs TCP sockets and OS threads, neither of
# which exist under WebAssembly. Pyodide builds disable default features.
# The server plans SQL, so it also needs DataFusion's default features
# (SQL, function libraries, Parquet, compression).
flight = ["dep:arrow-flight", "dep:prost", "dep:tonic", "datafusion/default", "tokio/net", "tokio/rt-multi-thread"]


# Pyodide downloads the wheel on every page load, so wasm builds trade
# speed for size: `MATURIN_PEP517_ARGS="--no-default-features --profile wasm"`.
[profile.wasm]
inherits = "release"
opt-level = "z"
lto = true
codegen-units = 1
strip = true

[build-dependencies]
pyo3-build-config = "0.28"
Expand Down
7 changes: 4 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ import xarray_sql as xql
# 4x-daily surface air temperature on a lat/lon grid, 2013-2014.
ds = xr.tutorial.open_dataset('air_temperature')

ctx = xql.XarrayContext()
ctx = xql.Context()
ctx.from_dataset('air', ds, chunks=dict(time=100))

# A climatology — the mean annual cycle — computed in SQL: average air
Expand Down Expand Up @@ -110,8 +110,9 @@ ds = xr.open_zarr(
storage_options={'token': 'anon'} # Anonymous read from the public GCS bucket — no auth required.
)

ctx = xql.XarrayContext()
# Make sure to pass `chunks`!
ctx = xql.Context()
# `chunks` sets the partition size; without it, partitions follow the
# store's own chunks (one hour each here).
ctx.from_dataset('era5', ds, chunks=dict(time=6), table_names={
('time', 'latitude', 'longitude'): 'surface',
('time', 'level', 'latitude', 'longitude'): 'atmosphere',
Expand Down
47 changes: 41 additions & 6 deletions docs/engines.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ DataFusion is the built-in engine, wrapped in a session:
```python
import xarray_sql as xql

ctx = xql.XarrayContext()
ctx = xql.Context()
ctx.from_dataset("era5", ds, chunks={"time": 24})
result = ctx.sql("SELECT ... FROM era5").to_dataset()
```
Expand Down Expand Up @@ -342,14 +342,14 @@ stores `float32` as `float64` and `bool` as an integer, MySQL `bool` as
keep the result's type.

The cursor is a one-shot Arrow stream: `xql.to_dataset(cur, ...)`
round-trips eagerly, and `chunks=` needs `spill=True`.
reads it into memory, and `chunks=` needs `spill=True`.

## Serving over Flight SQL (no copy)

The adapters above bring an engine to the data. `xql.serve` does the
reverse and brings the data to remote clients: it starts an
[Arrow Flight SQL](https://arrow.apache.org/docs/format/FlightSql.html)
server over the same lazy tables `XarrayContext` uses.
server over the same lazy tables `Context` uses.

```python
import xarray_sql as xql
Expand Down Expand Up @@ -453,7 +453,7 @@ to queries that filter on time.
TABLE`, `COPY`, `SET`, ...) are rejected, so clients cannot read or
write the server's filesystem.
- **DataFusion's SQL, without xarray-sql's Python UDFs.** The
`cftime()` and `reproject()` functions `XarrayContext` registers are
`cftime()` and `reproject()` functions `Context` registers are
not available on the server.
- **Bound its memory.** Any client can send an expensive `ORDER BY`,
join, or aggregation. `memory_limit=` (bytes) caps what those hold at
Expand All @@ -475,7 +475,7 @@ What each integration provides. Known issues and constraints live on

| | DataFusion | DuckDB | Polars | ADBC |
|---|---|---|---|---|
| Register | `XarrayContext` / any `SessionContext` | `xql.register(con, name, ds)` | `pl.scan_pyarrow_dataset(xql.arrow_dataset(ds))` | `xql.register(con, name, ds)` (copies into the database) |
| Register | `Context` / any `SessionContext` | `xql.register(con, name, ds)` | `pl.scan_pyarrow_dataset(xql.arrow_dataset(ds))` | `xql.register(con, name, ds)` (copies into the database) |
| Projection pushdown | yes | yes | yes | n/a (the database's own tables) |
| Chunk pruning on dim predicates | yes | yes | yes | n/a (the database's own indexes) |
| Eager round-trip (`xql.to_dataset`) | yes | yes | yes | yes (pass the cursor) |
Expand All @@ -501,7 +501,9 @@ only the source chunks it maps onto.
```mermaid
flowchart TB
R["xql.to_dataset(result, ...)"] --> K{"chunks=?"}
K -- "None (default)" --> E["eager: materialize once<br/>max_result_bytes= guards both the<br/>Arrow stream and the dense grid"]
K -- "None (default)" --> LZ{"result type"}
LZ -- "re-executable<br/>(DuckDB, Polars, DataFusion)" --> LI["lazily indexed, no dask:<br/>each access re-runs the query<br/>for its selection; .load() runs it<br/>once per variable (small results:<br/>read in one pass)"]
LZ -- "one-shot Arrow stream,<br/>or max_result_bytes=" --> E["in memory: read once<br/>max_result_bytes= guards both the<br/>Arrow stream and the dense grid"]
K -- "mapping / auto / inherit" --> SP{"spill=?"}
SP -- "False (default)" --> HD{"result type"}
HD -- "Polars LazyFrame/DataFrame<br/>DataFusion DataFrame" --> RX["re-execution: each window<br/>re-runs the query narrowed to its<br/>coordinate range (flows back into<br/>chunk pruning at the source)"]
Expand All @@ -515,6 +517,39 @@ few windows of a huge result. Spill pays one full pass plus temporary
disk — right when you'll touch most of the result, when the producer
is a DuckDB relation, or when all you have is a one-shot stream.

**Without chunks:** `chunks=None`, the default, returns lazily indexed
variables, as `xr.open_dataset(chunks=None)` does. Each access, such as
`out.t2m.sel(time="2020-01").values`, re-executes the query narrowed to
that selection and projected to the variable it reads, so it reads only
the source chunks under it, and `.load()` executes the query once per
variable. Nothing is chunked, so dask is not needed, and DuckDB
relations work too: every access runs on the calling thread, unlike
chunked windows, which dask computes on worker threads.

```python
out = xql.to_dataset(con.sql("SELECT * FROM era5"), template=ds)
out.t2m.sel(time="2020-01-01").mean().item() # reads one day's chunks
```

A lazy result must know its coordinates before any data is read, so
unless `coords="template"` supplies them, the query is first streamed
once. A small result (up to 64 MiB) is kept in memory from that pass,
so an aggregation, which no window filter can push below, executes
once. A larger one keeps only each dimension's distinct values and
stays lazy. One-shot Arrow streams, which cannot re-execute, and calls
with `max_result_bytes=` are read into memory directly.

Because a lazy result re-executes its query on access:

- the connection must stay open while the result is used: after
`con.close()`, reading a variable raises;
- changes to the source tables after `to_dataset` show through in
later reads;
- a query whose rows vary between runs, such as `LIMIT` without
`ORDER BY` or random sampling, can disagree with the coordinates
discovered for it. Read such a result once, in memory, with
`max_result_bytes=`.

Two knobs matter at scale:

- `coords="template"` trusts the template's coordinate arrays instead of
Expand Down
18 changes: 9 additions & 9 deletions docs/examples.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ import xarray_sql as xql

ds = xr.tutorial.open_dataset('air_temperature')

ctx = xql.XarrayContext()
ctx = xql.Context()
ctx.from_dataset('air', ds, chunks=dict(time=100))

clim = ctx.sql('''
Expand Down Expand Up @@ -66,17 +66,17 @@ url = 'gs://gcp-public-data-arco-era5/ar/full_37-1h-0p25deg-chunk-1.zarr-v3'
full = xr.open_zarr(url, chunks=None, storage_options={'token': 'anon'})

# A full year of hourly ERA5 — all 273 variables. No spatial slicing on the
# xarray side; SQL WHERE clauses below express the filters. `chunks={'time': 1}`
# aligns Dask chunks to native Zarr chunks of shape (1, 37, 721, 1440) so
# chunk reads from GCS happen concurrently.
# xarray side; SQL WHERE clauses below express the filters. Each table
# partition is one native Zarr chunk of shape (1, 37, 721, 1440); pass
# `chunks=` to `from_dataset` to group several into one partition.
#
# Heads up: 262 of those variables are surface and 11 are atmospheric. The
# library pushes column projection down, so SELECT only fetches what you ask
# for — but `SELECT * FROM era5.surface` would try to pull every variable
# across the year (terabytes from GCS). Always SELECT specific columns.
ds = full.sel(time='2020').chunk({'time': 1})
ds = full.sel(time='2020')

ctx = xql.XarrayContext()
ctx = xql.Context()
ctx.from_dataset('era5', ds, table_names={
('time', 'latitude', 'longitude'): 'surface',
('time', 'level', 'latitude', 'longitude'): 'atmosphere',
Expand Down Expand Up @@ -123,7 +123,7 @@ scalars into a single one-row table named `scalar`:
```python
import fsspec
import xarray as xr
from xarray_sql import XarrayContext
from xarray_sql import Context

# A real GOES-16 ABI cloud-and-moisture file from NOAA's public bucket:
# (y, x) image bands alongside dozens of scalar metadata variables.
Expand All @@ -135,7 +135,7 @@ ds = xr.open_dataset(fsspec.open_local(f'simplecache::{url}')).chunk(
{'y': 250, 'x': 250}
)

ctx = XarrayContext()
ctx = Context()
ctx.from_dataset('goes', ds)

# The gridded bands and the scalar metadata are separate tables.
Expand All @@ -153,7 +153,7 @@ A runnable version of the ERA5 example lives at

## The same tables on DuckDB and Polars

Every example above registers through an `XarrayContext`, but the tables are
Every example above registers through a `Context`, but the tables are
not DataFusion-specific: `xql.register(con, name, ds)` attaches the same lazy,
pushdown-scanned table to a DuckDB connection, and
`pl.scan_pyarrow_dataset(xql.arrow_dataset(ds))` serves Polars — same
Expand Down
Loading
Loading