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)
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):
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.
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