This tutorial builds a small analytics pipeline over two related datasets — employees and the departments they belong to. Along the way you’ll use array construction, compute kernels, row filtering, sorting and a hash join. Every cell runs when the docs are built.
The data
We start with two RecordBatches. The first holds employees; the second maps department ids to names.
Filtering rows has two spellings, and it is worth seeing both.
The query layer is the one you will reach for. col(...) names a column, comparison operators build a predicate, and the engine applies it to every column for you:
Note the comparison needs both operands to share a dtype — hence the array of thresholds. The query layer handles that for you, which is most of why it is the shorter spelling.
Step 2 — sort
sort_by orders a batch by one or more keys. Pass a list of (column, "ascending" | "descending") tuples; the second argument controls null placement (None uses the default). Multi-column keys are honoured — ties on the first key break on the next:
join performs a hash join between two record batches on named key columns. The signature is join(right, keys, right_keys=None, join_type="inner", num_threads=0):
keys — list of left-side key column names,
right_keys — right-side names, or None to reuse keys,
join_type — "inner", "left", "right", "full", "semi" or "anti",
num_threads — 0 auto-selects cores, 1 forces the serial path.
Only the join is eager here — it is a RecordBatch method. Everything after it is one plan, which runs when collect() is called. Add .optimize() before collect() to let the plan rewriter act on it first.
Step 5 — aggregate
Aggregation is not an eager kernel call — it happens inside a query plan. Wrap the batch with ma.memtable() to get a lazy table, describe what you want, and call .collect() to run it. Nulls are skipped: