On this page
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.
| Repositories | apache/datafusion, apache/arrow-rs |
| License | Apache-2.0 (Apache Software Foundation projects) |
| Versions | DataFusion 55.1 and arrow-rs 60.0 are current; Loam pins DataFusion 54.1 and arrow 58.4 (below) |
| In Loam | Available |
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.
- Planning. SQL (through
sqlparser-rs) or a DataFrame builds a logical plan: a tree of relational operators over table providers. - Optimizing. Rule-based rewrites push filters and projections toward scans, simplify expressions, reorder joins and turn
ORDER BY … LIMITinto top-k. - Physical planning. The logical plan becomes a tree of
ExecutionPlanoperators, each producing a number of output partitions. - 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
MemoryPoolso 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
| Addition | What it is |
|---|---|
| Table providers | Collections as tables, with a CollectionProvider that pushes down only the filters it can evaluate exactly on indexed fields |
| Retrieval operators | AnnExec (vector top-k), TantivySearchExec (BM25 top-k with global statistics), SparseExec, FilterBitmapExec |
| Combination operators | FusionExec (RRF, weighted, distribution-based), TailMergeExec (durable results ∪ log tail with upsert semantics), DocFetchExec |
| Table functions | vector_search, text_search, hybrid_search and rrf, with the same names Spice uses, so hybrid retrieval is reachable from plain SQL |
| Flight SQL server | Built 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.