Lazy queries

Build a plan, run it once, in Python

Alongside the eager PyArrow-shaped API, marrow has a lazy relational frontend: you describe a query, and nothing executes until collect().

The two worlds are told apart by the entry point — lazy is memtable / read_parquet (ibis spellings, which the frontend is modelled on); eager is record_batch / table (PyArrow spellings).

batch = ma.record_batch({
    "region": ma.array(["east", "west", "east", "west", "east"]),
    "product": ma.array(["a", "b", "a", "c", "b"]),
    "qty": ma.array([1, 2, None, 4, 5]),
    "price": ma.array([10, 20, 30, 40, 50]),
})

t = ma.memtable(batch)          # lazy, over an in-memory batch
t
LazyTable
InMemoryTable(5 rows)

Every verb returns a new plan. The underlying description is immutable and copying it is a refcount bump, so a plan is a reusable template rather than a builder you consume.

A first query

top = (
    t.filter(col("price") > lit(15))
     .aggregate(by=["region"], total=("sum", "price"))
     .order_by(("total", "descending"))
)
top.collect().to_pylist()
[{'region': 'east', 'total': 80}, {'region': 'west', 'total': 60}]

Nothing ran until .collect(). Before that, explain() shows the plan:

print(top.explain())
Sort(Aggregate(Filter(InMemoryTable(5 rows), gt(price, 15)), by=region, sum(price)), total desc)

The verbs

t.select("region", "price")            # project by name
t.drop("qty")
t.rename({"price": "amount"})
t.with_columns(total=col("qty") * col("price"))   # add/replace, keep the rest
t.project(only=col("price") + lit(1))  # replace the projection entirely
t.filter(col("region") == lit("east"))
t.aggregate(by=["region"], n=ma.count_star())
t.order_by(("price", "descending")).limit(10, offset=100)
t.head(5)
t.join(other, on="region", how="inner")
t.explain()                            # the plan as text
t.collect()  /  t.to_pyarrow()

collect(num_threads=) controls execution: 0 auto, 1 serial, N forces N workers.

Column expressions

col(name) reads a column, lit(value) is a constant, and the usual operators compose them:

expr = (col("price") * lit(2) > lit(60)) & col("qty").is_valid()
t.filter(expr).collect().to_pylist()
[{'region': 'west', 'product': 'c', 'qty': 4, 'price': 40},
 {'region': 'east', 'product': 'b', 'qty': 5, 'price': 50}]

Boolean logic is three-valued, like SQL. qty carries a null, so qty > 3 is null in that row — but null AND false is false, not null, because the second operand already decides the answer:

t.with_columns(
    big_qty=col("qty") > lit(3),          # null in the row where qty is null
    both=(col("qty") > lit(3)) & (col("price") > lit(100)),
).collect().to_pylist()
[{'region': 'east',
  'product': 'a',
  'qty': 1,
  'price': 10,
  'big_qty': False,
  'both': False},
 {'region': 'west',
  'product': 'b',
  'qty': 2,
  'price': 20,
  'big_qty': False,
  'both': False},
 {'region': 'east',
  'product': 'a',
  'qty': None,
  'price': 30,
  'big_qty': None,
  'both': False},
 {'region': 'west',
  'product': 'c',
  'qty': 4,
  'price': 40,
  'big_qty': True,
  'both': False},
 {'region': 'east',
  'product': 'b',
  'qty': 5,
  'price': 50,
  'big_qty': True,
  'both': False}]

The surface covers arithmetic and comparison operators, &/|/~/^, abs/sign/floor/ceil/round/trunc/sqrt/exp/ln, the string verbs (upper, lower, strip, length, startswith, contains, like, ilike, …), the temporal extractors and date_trunc, plus cast, isin, array_length, is_null, is_valid, is_nan, is_inf, fill_null, coalesce, nullif, if_else and case_when.

Aggregation

Keyword arguments name output columns; by= groups. An empty by is one implicit group:

t.aggregate(
    rows=ma.count_star(),
    values=("count", "qty"),
    total=("sum", "price"),
    avg=("mean", "price"),
).collect().to_pylist()
[{'rows': 5, 'values': 4, 'total': 150, 'avg': 30.0}]

count_star() counts rows and ("count", col) counts non-null values — they disagree on qty, which carries a null.

Functions: sum, mean, min, max, product, count, count_distinct, approx_count_distinct, variance, var_samp, stddev, stddev_samp.

HAVING needs no special support — filter the aggregate’s output:

(t.aggregate(by=["region"], total=("sum", "price"))
  .filter(col("total") > lit(50))
  .collect().to_pylist())
[{'region': 'east', 'total': 90}, {'region': 'west', 'total': 60}]

Group keys may be arbitrary expressions, and keys with no aggregates is SELECT DISTINCT:

t.aggregate(by=["region"]).collect().to_pylist()
[{'region': 'east'}, {'region': 'west'}]

Joins

Keys are given by name and resolved against each side’s schema:

managers = ma.memtable(ma.record_batch({
    "region": ma.array(["east", "west"]),
    "manager": ma.array(["ann", "bo"]),
}))

t.join(managers, on="region").select("region", "manager", "price").collect().to_pylist()
[{'region': 'east', 'manager': 'ann', 'price': 10},
 {'region': 'east', 'manager': 'ann', 'price': 30},
 {'region': 'east', 'manager': 'ann', 'price': 50},
 {'region': 'west', 'manager': 'bo', 'price': 20},
 {'region': 'west', 'manager': 'bo', 'price': 40}]

how= takes inner, left, right, full, semi and anti. Use left_on= / right_on= when the key names differ.

Reading Parquet

read_parquet opens the file’s footer metadata only — no column data is touched until collect():

t = ma.read_parquet("hits.parquet")
(t.filter(col("AdvEngineID") != lit(0))
  .aggregate(by=["RegionID"], users=("count_distinct", "UserID"))
  .order_by(("users", "descending"))
  .head(10)
  .collect())
WarningParquet BYTE_ARRAY columns arrive as binary

A column with no logical type is binary, and BinaryType is not a StringLikeType, so the string verbs reject it. Cast first:

col("url").cast(ma.string()).like("%google%")

Getting results out

result = t.filter(col("price") > lit(20))
print(type(result.collect()))     # a RecordBatch
print(result.to_pyarrow())        # zero-copy, over the C Data Interface
<class 'marrow.tabular.RecordBatch'>
pyarrow.RecordBatch
region: string
product: string
qty: int64
price: int64
----
region: ["east","west","east"]
product: ["a","c","b"]
qty: [null,4,5]
price: [30,40,50]

Limits worth knowing

  • read_parquet reads one file (local, or an object-store URI through OpenDAL): no globs, directories or hive partitions, and no CSV or JSON.
  • Row groups are skipped only after .optimize(). Without it, a filtered scan reads every row group.
  • No distinct/union/except/intersect nodes, and no regex.
  • Comparisons follow SQL, not IEEE: nan = nan is true here and false in marrow.compute — see Compute.

The full list is on Status & limitations.

Back to top