Tables & joins

A RecordBatch is a schema paired with a list of equal-length column arrays — the fundamental tabular unit in Arrow. A Table is a schema paired with chunked columns, so it can span many batches without copying them into one contiguous block. Both mirror their PyArrow namesakes.

Creating a RecordBatch

The most direct constructor takes a dict mapping column names to arrays:

batch = ma.record_batch({
    "id":    ma.array([1, 2, 3, 4]),
    "name":  ma.array(["Alice", "Bob", "Carol", "Dave"]),
    "score": ma.array([9.5, 8.2, 7.8, 9.1]),
})
print(batch)
RecordBatch(num_rows=4, schema=Schema(fields=[id: int64, name: string, score: float64]))

You can also pass a list of arrays with explicit names=:

batch2 = ma.record_batch(
    [ma.array([1, 2, 3]), ma.array(["x", "y", "z"])],
    names=["id", "label"],
)
print(batch2.column_names)
['id', 'label']

Inspecting a batch

print("shape:       ", batch.shape)        # (rows, columns)
print("num_rows:    ", batch.num_rows)
print("num_columns: ", batch.num_columns)
print("column_names:", batch.column_names)
print("schema:      ", batch.schema)
shape:        (4, 3)
num_rows:     4
num_columns:  3
column_names: ['id', 'name', 'score']
schema:       Schema(fields=[id: int64, name: string, score: float64])

Reach a single column by name or position — both return an array:

print(batch.column("name"))
print(batch.column(0))
StringArray([Alice, Bob, Carol, Dave])
PrimitiveArray[int64]([1, 2, 3, 4])

Convert the whole batch to native Python, either column-oriented (to_pydict) or row-oriented (to_pylist):

print(batch.to_pydict())
print(batch.to_pylist())
{'id': [1, 2, 3, 4], 'name': ['Alice', 'Bob', 'Carol', 'Dave'], 'score': [9.5, 8.2, 7.8, 9.1]}
[{'id': 1, 'name': 'Alice', 'score': 9.5}, {'id': 2, 'name': 'Bob', 'score': 8.2}, {'id': 3, 'name': 'Carol', 'score': 7.8}, {'id': 4, 'name': 'Dave', 'score': 9.1}]

Selecting, slicing, reshaping

select keeps a subset of columns (by name or index); slice takes a zero-copy window of rows:

print(batch.select(["name", "score"]).to_pydict())
print(batch.slice(1, 2).to_pylist())         # offset 1, length 2
{'name': ['Alice', 'Bob', 'Carol', 'Dave'], 'score': [9.5, 8.2, 7.8, 9.1]}
[{'id': 2, 'name': 'Bob', 'score': 8.2}, {'id': 3, 'name': 'Carol', 'score': 7.8}]

Column mutations return a new batch — arrays are immutable. Each takes a Field describing the new column, built with ma.field(name, type, nullable, metadata):

rank = ma.field("rank", ma.int64(), True, None)

with_rank     = batch.append_column(rank, ma.array([1, 2, 3, 4]))
inserted      = batch.add_column(0, rank, ma.array([1, 2, 3, 4]))
replaced      = batch.set_column(2, rank, ma.array([1, 2, 3, 4]))
renamed       = batch.rename_columns(["ID", "NAME", "SCORE"])
without_score = batch.remove_column(2)

print("appended: ", with_rank.column_names)
print("inserted: ", inserted.column_names)
print("replaced: ", replaced.column_names)
print("renamed:  ", renamed.column_names)
print("removed:  ", without_score.column_names)
appended:  ['id', 'name', 'score', 'rank']
inserted:  ['rank', 'id', 'name', 'score']
replaced:  ['id', 'name', 'rank']
renamed:   ['ID', 'NAME', 'SCORE']
removed:   ['id', 'name']

Sorting

sort_by orders a batch by a column. The second argument is null placement (None for the default, or "at_start" / "at_end"):

print(batch.sort_by([("score", "descending")]).to_pylist())
[{'id': 1, 'name': 'Alice', 'score': 9.5}, {'id': 4, 'name': 'Dave', 'score': 9.1}, {'id': 2, 'name': 'Bob', 'score': 8.2}, {'id': 3, 'name': 'Carol', 'score': 7.8}]

Multi-column keys work too — pass more (column, direction) tuples and ties on the first key break on the next. See Sorting for the full story.

Joins

join runs a hash join between two batches on named key columns:

people = ma.record_batch({
    "dept_id": ma.array([1, 2, 1, 3]),
    "name":    ma.array(["Alice", "Bob", "Carol", "Dave"]),
})
depts = ma.record_batch({
    "dept_id": ma.array([1, 2, 3]),
    "dept":    ma.array(["Engineering", "Sales", "Marketing"]),
})

inner = people.join(depts, ["dept_id"], join_type="inner")
print(inner.to_pylist())
[{'dept_id': 1, 'name': 'Alice', 'dept_id_right': 1, 'dept': 'Engineering'}, {'dept_id': 1, 'name': 'Carol', 'dept_id_right': 1, 'dept': 'Engineering'}, {'dept_id': 2, 'name': 'Bob', 'dept_id_right': 2, 'dept': 'Sales'}, {'dept_id': 3, 'name': 'Dave', 'dept_id_right': 3, 'dept': 'Marketing'}]

The signature is join(right, keys, right_keys=None, join_type="inner", num_threads=0):

Argument Meaning
right the right-hand RecordBatch
keys list of left-side key column names (a bare string is also accepted)
right_keys right-side names, or None to reuse keys
join_type "inner", "left", "right", "full", "semi", "anti"
num_threads 0 auto-selects cores; 1 forces the serial path

The supported join types:

for kind in ["inner", "left", "right", "full", "semi", "anti"]:
    out = people.join(depts, ["dept_id"], join_type=kind)
    print(f"{kind:>6}: {out.num_rows} rows, columns={out.column_names}")
 inner: 4 rows, columns=['dept_id', 'name', 'dept_id_right', 'dept']
  left: 4 rows, columns=['dept_id', 'name', 'dept_id_right', 'dept']
 right: 4 rows, columns=['dept_id', 'name', 'dept_id_right', 'dept']
  full: 4 rows, columns=['dept_id', 'name', 'dept_id_right', 'dept']
  semi: 4 rows, columns=['dept_id', 'name']
  anti: 0 rows, columns=['dept_id', 'name']

Columns whose names collide across the two inputs are suffixed with _right, the right-hand key column included. PyArrow’s defaults differ: it merges the key columns and adds no suffix.

Grouping & aggregation

Aggregation lives in the query layer, not on RecordBatch. Wrap a batch with ma.memtable() to get a lazy table, describe the reduction, and .collect() runs it and hands back a RecordBatch:

sales = ma.record_batch({
    "region": ma.array(["EU", "US", "EU", "US", "EU"]),
    "amount": ma.array([10, 20, 30, 40, 50]),
})

q = ma.memtable(sales)
print(q.aggregate(
    total=("sum", "amount"),
    average=("mean", "amount"),
    n=("count", "amount"),
).collect().to_pylist())
[{'total': 150, 'average': 30.0, 'n': 5}]

A by= key turns the same call into a GROUP BY — one row per group. Keyword arguments name the output columns:

by_region = q.aggregate(
    by=["region"],
    total=("sum", "amount"),
    distinct=("count_distinct", "amount"),
).order_by("region")
print(by_region.collect().to_pylist())
[{'region': 'EU', 'total': 90, 'distinct': 3}, {'region': 'US', 'total': 60, 'distinct': 2}]

Available functions: count, sum, product, mean, min, max, count_distinct, approx_count_distinct (HyperLogLog), variance, var_samp, stddev and stddev_samp. ma.count_star() is COUNT(*), which differs from ("count", col) on a nullable column. Group keys may be arbitrary expressions, not just names, and keys with no aggregates is SELECT DISTINCT.

HAVING needs no special support — .aggregate(...).filter(...) evaluates against the aggregate’s own output.

Tables

A Table holds chunked columns and is built the same way as a batch, from a dict or from arrays with names=:

table = ma.table({
    "id":   ma.array([1, 2, 3, 4, 5]),
    "city": ma.array(["NYC", "LA", "SF", "CHI", "SEA"]),
})
print(table)
print("shape:", table.shape)
Table(num_rows=5, num_columns=2, schema=Schema(fields=[id: int64, city: string]))
shape: (5, 2)

Tables share the inspection API with batches — schema, column_names, column, to_pydict, to_pylist — and expose their underlying batches:

print(table.column("city"))
print(table.to_pydict())
print("batches:", len(table.to_batches()))
ChunkedArray([StringArray([NYC, LA, SF, CHI, SEA])])
{'id': [1, 2, 3, 4, 5], 'city': ['NYC', 'LA', 'SF', 'CHI', 'SEA']}
batches: 1

A table can carry many batches at once, each kept as its own chunk, so building a table from batches does not copy them. Most Table verbs combine the chunks before they run, and all execution is in memory.

Note

Table.sort_by and Table.join combine the table’s chunks into one batch before they run. Aggregation is on neither type; it is a query-layer verb, as above.

Back to top