Skip to content

ArrowMetal 0.4.0

Latest

Choose a tag to compare

@singhpratech singhpratech released this 02 Oct 17:58
· 26 commits to main since this release

Everything below is new in 0.4.0. datafusion-arrowmetal, on crates.io, is a physical optimizer rule
for Apache DataFusion 55.1: registered on a SessionContext, it runs full sorts of 250,000 rows and
more on the GPU, 6.9x to 28.8x faster than DataFusion alone up to 50,000,000 rows
(datafusion/results/datafusion_sort_warm_2026-09-29.csv), and the count(*), DISTINCT and integer
MIN/MAX group-bys its measured table takes over MemTables of 50,000,000 rows or more, 2.31x to
4.27x faster warm (datafusion/results/datafusion_groupby_refit_check_2026-10-02.csv,
datafusion_groupby_default_import_2026-10-01.csv); 13,632 query pairs run with and without the rule
give the same answers. The Polars MetalEngine passes every Polars sort key as one key with Polars'
null placement and float order, so a sort keeps Polars' key count and sort().head(), top_k and
bottom_k run as the GPU top-k (the top 100 by a nullable Float64 key descending over 50,000,000 rows,
shapes="all": 129.53 → 11.07 ms, Benchmarks/results/polars_engine_sort_options_2026-09-30.csv); its
group-by crossovers are refitted and held to the default benchmark, where the default takes 54 of the
150 group-by case-size pairs, each 1.12x to 10.69x of the faster Polars engine on the best run
(Benchmarks/results/polars_engine_default_groupby_2026-09-30.csv); the conformance grid compares
15,376 cases with Polars, 0 unclassified. Every query the DuckDB rewrite's auto mode rewrites is 1.09x
to 5.03x faster than DuckDB on the best run (Benchmarks/results/duckdb_rewrite_2026-10-02.csv). A
column held as many Arrow arrays imports in one call with no concatenated copy (C, Swift, Python, Rust,
Go, Node, R), and sorts take a null placement and a float order per key (ieee, total,
nan_largest). Correctness: kernels over arrays near 2^32 elements, and row-wise grids of 2^32 threads,
no longer wrap and skip rows; Parquet's fixed-width decode no longer overwrites a column whose decoded
bytes pass 4 GB; count(expr) in a plan's group_by and the streaming group-by's count(column) count
the non-null values for every type; the streaming external sort merges NaN where the GPU sort placed it;
Node imports empty typed arrays and empty Arrow vectors. Behaviour changes: index arrays (argsort,
top-k, lexsort, the ranks, join indices, a group's representative rows) are UInt32, and rows past
2^32 - 1 are refused; a Float64 group sum is the correctly rounded sum of the group's values (equal
to math.fsum, independent of row order) and a Float64 group mean that sum divided by the count,
rounded once; the Polars engine's sort_helper_keys shape class is gone, and
MetalEngine(shapes={"sort_helper_keys"}) raises ValueError.

The engines, measured

Apple M4 Max (16 CPU cores, 64 GB). Each speed-up is the engine's own time divided by the time with
ArrowMetal, with the CPU time of the call in brackets. Every figure is from the results file named
beside it.

DataFusion (datafusion-arrowmetal, the default configuration):

rows DataFusion alone, ms (CPU-ms) with the rule, ms (CPU-ms) speed-up
ORDER BY an int64 key, 3 columns 50,000,000 1,581.22 (6,104.1) 69.74 (125.8) 22.7x
count(*) over two int32 keys, 1,000,000 groups 50,000,000 75.36 (1,127.5) 17.63 (61.7) 4.27x
DISTINCT over two int32 keys, 200 groups 50,000,000 30.16 (422.5) 13.05 (47.7) 2.31x

Full sorts of 250,000 to 50,000,000 rows: 6.9x to 28.8x
(datafusion/results/datafusion_sort_warm_2026-09-29.csv). The ten group-by series the default runs
on the GPU (20 cases, 50,000,000 rows): 2.31x to 4.27x warm; on the first run after 500 ms of idle
against DataFusion's first run after the same idle, 1.24x to 2.17x, and after 5 s, 1.11x to 1.68x
(datafusion/results/datafusion_groupby_refit_check_2026-10-02.csv,
datafusion/results/datafusion_groupby_default_import_2026-10-01.csv). To improve, left to DataFusion
by the default: top-k (ORDER BY … LIMIT 100) 0.14x to 0.41x
(datafusion/results/datafusion_sort_warm_2026-09-29.csv), filters 0.31x to 0.79x
(datafusion/results/datafusion_filter_2026-10-01.csv).

Polars (MetalEngine()):

50,000,000 rows faster Polars engine, ms (CPU-ms) MetalEngine(), ms (CPU-ms) speed-up
group-by over 2 keys, 1,000,000 groups, count streaming 224.9 (3,207.1) 23.9 (15.8) 9.42x
sort 3 columns by an int64 key in-memory 377.8 (4,647.8) 74.6 (27.8) 5.07x
inner join, 1,000,000-row build side streaming 82.2 (1,223.7) 35.2 (13.0) 2.33x
group-by over 1 key, 100,000 groups, mean streaming 60.2 (741.7) 49.9 (14.6) 1.21x

(Benchmarks/results/polars_engine_bench_2026-09-26-final3.csv.) Group-bys under the refitted table:
54 of 150 case-size pairs taken, 1.12x to 10.69x of the faster Polars engine on the best run
(Benchmarks/results/polars_engine_default_groupby_2026-09-30.csv). Conformance: 15,376 cases,
15,096 identical, 33 within the float-summation bound, 0 different, 247 where Polars' plan had nothing
to run (Benchmarks/results/engine_conformance_2026-10-02.csv). To improve, left to Polars by the
default: whole-frame aggregates 0.09x to 0.24x, row-wise filters and projections 0.14x to 0.42x, a
semi join against a 1,000-row table 0.21x to 0.47x, the top 100 by a nullable Float64 key 0.75x to
1.05x (Benchmarks/results/polars_engine_crossover_2026-09-30-groupby.csv).

DuckDB (the rewrite extension, auto):

rows DuckDB, ms (CPU-ms) rewritten, ms (CPU-ms) speed-up
100,000 INTEGER keys: sum, count 50,000,000 79.44 (1,199.5) 15.79 (100.5) 5.03x
sum, max, avg (BIGINT), no GROUP BY 50,000,000 4.33 (62.8) 3.98 (49.3) 1.09x

Every query auto rewrote: 1.09x to 5.03x on the best run, 1.04x to 4.66x on the median
(Benchmarks/results/duckdb_rewrite_2026-10-02.csv). Conformance: 33,376 generated queries, 21,844
rewritten and identical, 11,532 left to DuckDB, 0 different
(Benchmarks/results/engine_conformance_2026-09-25.csv). To improve: on the first run after 500 ms of
idle, five of the thirteen rewritten query-size pairs are behind DuckDB's first run after the same idle
(0.47x to 0.99x, Benchmarks/results/duckdb_rewrite_2026-10-02.csv).

Install

pip install arrowmetal                      # macOS 14+ on Apple silicon; the wheel is attached below
pip install "arrowmetal[polars]"            # with Polars, in the tested range
# Cargo.toml
arrowmetal = "0.4.0"                        # the arrow-rs binding
datafusion-arrowmetal = "0.4.0"             # the DataFusion rule, with datafusion = "=55.1.0"
.package(url: "https://github.com/singhpratech/ArrowMetal", from: "0.4.0")
go get github.com/singhpratech/ArrowMetal/go/arrowmetal@v0.4.0

The Rust crates and the Go module link libArrowMetalC.dylib from a Swift build of the repository
or from the wheel (site-packages/arrowmetal/_lib/libArrowMetalC.dylib), as docs/RUST.md and
docs/GO.md describe. The DuckDB extensions are built from the repository against DuckDB 1.5.5
(duckdb-extension/build.sh, duckdb-extension/build_rewrite.sh). Every change is listed in
CHANGELOG.md, section 0.4.0.