Queries in Mojo
One plan IR, two lanes — comptime-typed and runtime-typed
The Python frontends on the previous pages are drivers for an engine written in Mojo. This page is that engine’s own API — and the thing that makes it unusual: one relational plan IR with two expression lanes underneath it.
| Comptime lane | Runtime lane | |
|---|---|---|
| You write | col("a", int64) |
col("a") |
| Operand dtype resolved | at compile time | at run time |
| A subtree becomes | one fused SIMD loop, no dispatch | an interpreted tree |
| Costs you | compile time | a tag switch per node |
| Lives in | marrow/expr/comptime/ |
marrow/expr/runtime/ |
Both erase into the same box — DynValue — so every relational operator is compiled exactly once, and a plan can mix lanes at the box, never inside a node. Picking a lane is the only decision; the verbs are identical. What the comptime lane buys in speed is measured query by query.
The verbs
A plan is a sentence, read left to right:
var out = (
table(orders)
.filter(col("amount", int64) > lit(5, int64))
.limit(2)
.execute()
).filter(), .select(), .project(), .with_columns(), .drop(), .rename(), .aggregate(), .sort_by(), .limit() and .join() chain off any relation. table(batch) and scan(path, schema) are the two sources, and .execute() drains the plan into a single RecordBatch.
Aggregates are methods on the column
Here is the whole thing as a complete, compiling program — filter, group, HAVING, order, limit:
docs/snippets/query_pipeline.mojo
"""A whole query as one sentence — filter, group, having, order, limit.
Compiled by `pixi run -e dev docs_check`; included by docs/guide/expressions.qmd.
"""
from marrow.builders import array
from marrow.dtypes import int64
from marrow.expr import col, lit, table
from marrow.tabular import record_batch
def main() raises:
var orders = record_batch(
[
array([1, 2, 1, 2], int64).copy(),
array([10, 20, None, 40], int64).copy(),
],
names=["customer", "amount"],
)
var out = (
table(orders^)
.filter(col("amount", int64) > lit(5, int64))
.aggregate(
[col("amount", int64).sum().alias("total")],
[col("customer", int64)],
)
.filter(col("total", int64) > lit(5, int64))
.sort_by([col("total", int64)], [False])
.limit(1)
.execute()
)
print(out)That file is compiled by CI, not transcribed into the page — pixi run docs_check builds every snippet under docs/snippets/, so this listing cannot drift from an API that has moved.
Aggregates come first and keys second — aggregate(aggs, keys) — and an omitted key list is a whole-table reduction. sum, product, mean, min, max, count, count_distinct, approx_count_distinct, variance and stddev are all methods, named with .alias(...). count_star() is COUNT(*).
Never spell the comptime Aggregate[Fold[SumKernel, Int64Type], ...] node by hand. That form is for the kernel layer and the binary-size benches; the fluent method is the API.
Window functions
Attach a window to any aggregate or ranking function with .over(...):
docs/snippets/window_functions.mojo
"""Window functions — ranking and a partitioned running total.
Compiled by `pixi run docs_check`; included by docs/guide/expressions.qmd.
"""
from marrow.builders import array
from marrow.dtypes import int64, string
from marrow.expr import col, row_number, table
from marrow.tabular import record_batch
def main() raises:
var sales = record_batch(
[
array(["east", "west", "east", "west"]).to_dyn(),
array([10, 20, 30, 40], int64).to_dyn(),
],
names=["region", "amount"],
)
var out = (
table(sales^)
.with_columns(
["rn", "running"],
[
row_number().over(order_by=[col("amount", int64)]),
col("amount", int64)
.sum()
.over(
partition_by=[col("region", string)],
order_by=[col("amount", int64)],
),
],
)
.execute()
)
print(out)row_number, rank and dense_rank are re-exported from marrow.expr; percent_rank, cume_dist and ntile come from marrow.expr.builders. lag, lead, first_value, nth_value and last_value are methods on a value. A window preserves input order, not sorted order.
The optimizer
execute() alone optimizes nothing. The rewriter is a comptime parameter, so a binary links exactly the rules it names:
from marrow.expr import AllRules
var optimized = plan.optimize[AllRules]()
print(optimized) # the rewritten plan, as text
var out = optimized.execute()AllRules is sixteen rules — elimination, merging, conjunction splitting, six pushdowns (five for filters, one for limits) and a Top-N rewrite — plus ColumnPruning, a preparatory downward pass. They run in a chosen order, so each sees the previous one’s output. Because a plan renders as text, optimize[NoRules]() and optimize[AllRules]() can be printed and diffed.
Late-bound parameters
param() sits beside col() and lit(): col reads from data, lit is a constant, param is a constant supplied later.
docs/snippets/late_bound_params.mojo
"""A late-bound parameter, supplied per execution through `Bindings`.
Compiled by `pixi run docs_check`; included by docs/guide/expressions.qmd.
"""
from marrow.builders import array
from marrow.dtypes import int64
from marrow.expr import col, param, table
from marrow.scalars import Int64Scalar
from marrow.tabular import record_batch
def main() raises:
var batch = record_batch(
[array([1, 5, 9], int64).to_dyn()], names=["a"]
)
var min_a = param("min-a", int64)
var plan = table(batch^).filter(col("a", int64) > min_a.copy())
# One plan, two executions -- the value travels through, not into, the plan.
print(plan.execute(bindings={"min-a": Int64Scalar(4).to_dyn()}))
print(plan.execute(bindings={"min-a": Int64Scalar(8).to_dyn()}))The value travels through the execution as a Bindings map rather than being substituted into a copy of the plan, so two executions of one plan cannot interfere — and the fused inner loop is unchanged, so a parameter costs nothing per row. Give it a default= to make it optional.
This is what compiling a query is built on.
The push engine
Relation.to_operator(ctx) lowers a plan node into an Operator that owns all the mutable state. The engine pushes: push(batch) answers with what it produced, drain() with what is left, and collect() runs a plan down to one StructArray. There is one operator per plan node.
Choosing a lane
Use the comptime lane (col("a", int64)) when the schema is known where the code is written and you want the fusion and the small binary — this is the lane marrow compile exists for.
Use the runtime lane (col("a")) when the query is assembled from strings at run time. It interprets by switching on a tag, which costs binary size only in programs that use it at all. That asymmetry is deliberate: a program that builds expressions at run time has already accepted an interpreter.
See Architecture for why the closed comptime world is dead-code-eliminable, and what that buys.