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)
Note

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(*).

Note

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.

Back to top