Loam is pre-alpha: the engine core runs today; Live and Durable are in progress. See the roadmap

Blog/What Loam is built on

Apache DataFusion and Arrow: the query engine under every Loam read

Every Loam read, whether a Qdrant query, an Elasticsearch search or SQL, compiles to a DataFusion plan over Arrow batches. How DataFusion works, what Loam adds to it, why we embed it, and the version lockstep that comes with it.

Apache DataFusion and Arrow: the query engine under every Loam read
On this page
  1. Arrow: the shared memory format
  2. DataFusion: how a query runs
  3. Why we embed it
  4. What Loam adds
  5. The version lockstep
  6. Limits

Loam speaks several protocols, but it has one query engine. A Qdrant query, an Elasticsearch _search, a native hybrid request and a SQL statement over Flight SQL all compile to an Apache DataFusion plan, and every intermediate result is an Apache Arrow record batch.

Repositoriesapache/datafusion, apache/arrow-rs
LicenseApache-2.0 (Apache Software Foundation projects)
VersionsDataFusion 55.1 and arrow-rs 60.0 are current; Loam pins DataFusion 54.1 and arrow 58.4 (below)
In LoamAvailable

Arrow: the shared memory format

Arrow is a specification for columnar data in memory. A record batch is a set of equal-length columns, each a contiguous buffer of values plus a validity bitmap for nulls, with offsets for variable-length types such as strings and lists. Because the layout is specified to the byte, a batch produced in Rust can be read in Python, C++ or Java without copying or converting it. Arrow also defines an IPC format for streaming batches between processes and Arrow Flight, a gRPC protocol for moving them over the network; Flight SQL layers SQL queries and metadata on top, and ADBC gives each language a database driver API over it.

For Loam, that means one data representation from the storage formats (Lance and Parquet are both Arrow-native) through the query operators to the wire. A Python client that asks for results over Flight SQL gets Arrow batches that PyArrow, Polars and pandas can use directly.

DataFusion: how a query runs

DataFusion is an extensible query engine written in Rust, distributed as a library.

  1. Planning. SQL (through sqlparser-rs) or a DataFrame builds a logical plan: a tree of relational operators over table providers.
  2. Optimizing. Rule-based rewrites push filters and projections toward scans, simplify expressions, reorder joins and turn ORDER BY … LIMIT into top-k.
  3. Physical planning. The logical plan becomes a tree of ExecutionPlan operators, each producing a number of output partitions.
  4. Execution. Each partition is an async stream of record batches, pulled by its parent, running on Tokio. Operators are vectorised over Arrow arrays, and memory is tracked in a MemoryPool so large operators can spill.

Every layer is extensible: custom table providers, logical and physical operators, optimizer rules, scalar, aggregate, window and table functions. That extensibility is why InfluxDB 3, GreptimeDB, LanceDB, DataFusion Comet, Sail and Spice all build on it.

Why we embed it

We did not want to write a query engine. We wanted to write the parts no engine has: vector and BM25 retrieval over Lance and Tantivy, fusion, and a merge with the unindexed log tail. DataFusion lets us add those as operators and reuse everything else: SQL, expression evaluation, joins, aggregations, sorting, memory accounting. The alternatives were C++ engines (DuckDB, Velox), which are excellent but harder to extend and embed from a Rust codebase, or writing our own, which is years of work that is not our differentiator. And Lance, our collection format, is itself built on DataFusion and Arrow.

What Loam adds

AdditionWhat it is
Table providersCollections as tables, with a CollectionProvider that pushes down only the filters it can evaluate exactly on indexed fields
Retrieval operatorsAnnExec (vector top-k), TantivySearchExec (BM25 top-k with global statistics), SparseExec, FilterBitmapExec
Combination operatorsFusionExec (RRF, weighted, distribution-based), TailMergeExec (durable results ∪ log tail with upsert semantics), DocFetchExec
Table functionsvector_search, text_search, hybrid_search and rrf, with the same names Spice uses, so hybrid retrieval is reachable from plain SQL
Flight SQL serverBuilt on arrow-rs's arrow-flight crate: queries, and DoPut bulk ingest into collections

A native hybrid request becomes: filter bitmap → (vector search ‖ BM25) → fusion → limit → fetch, as one plan with no network hop between retrieval and fusion. Row ids share one roaring-bitmap space across masks, fusion and fetch; durable rows carry their Lance stable row id, and tail documents carry ids above 263 so the two never collide. DataFusion 54's dynamic filters push top-k bounds from the fusion down into the scans.

For SQL, DataFusion accepts only positional arguments to table functions in FROM, so Loam's functions take them in a fixed order:

SELECT id, _score
FROM rrf(vector_search('docs', [0.1, 0.2, 0.3], 'embedding', 100),
         text_search('docs', 'object storage', 'body', 100))
LIMIT 10;

Planned: distributed execution with datafusion-distributed, which runs plan stages on several nodes and moves batches between them over Arrow Flight, for large scans, joins and aggregations; Iceberg scans with a hot tier; graph expansion as an ExpandExec operator. Planned

The version lockstep

This is the cost of building on DataFusion, and it is real. Every crate that extends DataFusion or Arrow has to use the same major versions as everything else in the binary, because their types cross crate boundaries. DataFusion releases a major version roughly monthly.

  • Lance pins the pair. Lance 12 depends on DataFusion 54 and arrow 58, so Loam is on 54.1 and 58.4 while upstream is on 55.1 and 60.0. We move when Lance moves.
  • It decides what we can link. Sail's crates are on DataFusion 55 and arrow 59, and Spice ships its own forks of DataFusion and arrow-rs. Linking either would mean two copies of the engine. So Loam integrates with them over protocols (Spark Connect, Flight SQL), never as libraries.
  • It decides what extensions we can use. datafusion-distributed 3.0 was the last release on DataFusion 54; datafusion-federation 0.5.5 likewise.
  • Flight SQL is marked experimental in arrow-rs (the feature is literally flight-sql-experimental), so its API can change with any arrow release.

We keep DataFusion's default features trimmed to what Lance needs, so adding a crate does not rebuild the engine with different features, and every bump is one pull request that moves Lance, DataFusion and Arrow together.

Limits

  • No distributed execution in the core project (Ballista and datafusion-distributed are separate, and younger).
  • Frequent breaking API changes between majors, the price of fast development.
  • The optimizer is rule-based with limited statistics, so plans for complex joins need care. Loam's retrieval queries are narrow enough that this has not mattered yet.

The next post covers the format that gives those operators their data: Lance.

More from the blog