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:
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())
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(" %g oogle%" )
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