Back to Blog

DataFusion on the Apple Silicon GPU: One Optimizer Rule, the Same SQL, Sorts 6.9x to 28.8x Faster

Prateek SinghOctober 2, 202610 min read11 views
DataFusion on the Apple Silicon GPU: One Optimizer Rule, the Same SQL, Sorts 6.9x to 28.8x Faster

Apache DataFusion runs any physical optimizer rule you register. ArrowMetal 0.4.0 ships one for the Apple silicon GPU: the SQL is unchanged, the answers are DataFusion's, and full sorts of 250,000 to 50 million rows run 6.9x to 28.8x faster on an M4 Max. What it takes, what it leaves, and why.

Apache DataFusion will run a physical optimizer rule you hand it, after its own. The rule sees the finished physical plan and may replace any node with an operator of its own; DataFusion runs what comes back and never looks at the SQL again. ArrowMetal 0.4.0, released on 2 October, is one such rule for the Apple silicon GPU, the Rust crate datafusion-arrowmetal. Register it on a SessionContext and EXPLAIN shows a MetalExec where SortExec was. The SQL is unchanged, the answer is DataFusion's, and a report says which nodes it took, which it left, and why.

optimizer rule added to DataFusion
1
the SQL and every other plan node stay DataFusion's
full sorts, 250,000 to 50,000,000 rows
6.9x to 28.8x
faster than DataFusion alone, four key types
query pairs, rule off against rule on
13,632
0 mismatches
CPU-ms to sort 50,000,000 rows
126 vs 6,104
with the rule, against 16 partitions without it

Apple M4 Max, DataFusion 55.1.0, ArrowMetal 0.4.0; warm, best of 5, 8,192-row batches, 29 September 2026; the grid on 1 October.

ArrowMetal is my Apache-2.0 Arrow compute library for the Apple silicon GPU, on Arrow's Powered By page. This post is about the DataFusion side.

A rule in DataFusion's own optimizer

DataFusion plans in two passes: logical rules, then physical optimizer rules over the operators that actually run, and that list is yours to extend. The crate's rule goes in before DataFusion's last two, the filter pushdown and the sanity check, so a rewritten plan is still checked. Two lines are new. This ran on my M4 Max from the crate's quickstart example; the plan and report are what it printed:

// datafusion = "=55.1.0", datafusion-arrowmetal = "0.4.1"
let rule = ArrowMetalRule::new(ArrowMetalConfig::default());   // full sorts from 250,000 rows
let ctx = session_context(SessionConfig::new(), rule.clone());
ctx.register_table("sales", Arc::new(MemTable::try_new(schema, vec![vec![batch]])?))?;

let plan = ctx.sql("SELECT name, region, amount FROM sales ORDER BY amount DESC NULLS LAST, name")
    .await?.create_physical_plan().await?;
println!("{}", displayable(plan.as_ref()).indent(false));
println!("{}", rule.report());

// 1,000,000-row MemTable; printed:
ProjectionExec: expr=[name@2 as name, region@0 as region, amount@1 as amount]
  MetalExec: sort=[amount DESC NULLS LAST, name ASC NULLS LAST]
    DataSourceExec: partitions=1, partition_sizes=[1]

TAKEN    SortExec: expr=[amount@1 DESC NULLS LAST, name@2 ASC NULLS LAST]
         -- input rows 1000000 (exact) vs min_rows 250000

The rule replaced the sort because every key is a column of a type it sorts and DataFusion's statistics gave the input an exact row count of at least 250,000. MemTables and DataFusion's own Parquet scan give exact counts; a sort above a filter or a join has an estimate, and the default leaves it. MetalExec emits one partition, so it takes the SortPreservingMergeExec DataFusion plans over 16 per-partition sorts too, and puts a RepartitionExec back when something higher needs one. The report also says no, with the reason, which is the part I use most:

-- ... ORDER BY amount DESC NULLS LAST LIMIT 3
LEFT     SortExec: TopK(fetch=3) -- top-k (sort with fetch 3) disabled in config

-- SELECT region, count(*), avg(amount) FROM sales GROUP BY region ORDER BY region
LEFT     AggregateExec: mode=FinalPartitioned -- input rows 1000000 (exact) vs min_rows 250000;
         count + sum_avg_f64 over 1 i64 key (memory input): the measured table takes it
         at no group count and size (results/datafusion_groupby_sweep_2026-10-01.csv)

What the default takes, and what it leaves

The take-list is a measured table, not a flag. Full sorts are taken from 250,000 rows because that is the first size at which every key type was ahead by 2.64x cold and 6.9x warm. The lead grows with the input: 20.3x to 26.1x at 2,000,000 rows, 19.0x to 26.6x at 50,000,000, where DataFusion alone takes 1.58 to 1.86 seconds and the rule 61 to 88 ms. Over DataFusion's own Parquet reader a Float64 sort of 10,000,000 and 50,000,000 rows is 10.8x to 16.5x, a file with a string column 3.6x to 5.5x, the decode counted inside the rule's time.

Full sorts in DataFusion, with the rule: speed-up over DataFusion alone ORDER BY one key, three columns carried · warm, best of 5 · 16 partitions, 8,192-row batches · rows on a log scale 1x5x10x15x20x25x30x 250K500K1M2M10M50M Float32 int64 String Float64
Full sorts, DataFusion alone ÷ with the rule, from datafusion_sort_warm_2026-09-29.csv: Apple M4 Max, 16 partitions, 8,192-row batches, warm, best of 5. One batch per partition gives 6.9x to 28.8x over the same cells.

Top-k is the honest opposite. ORDER BY … LIMIT 100 is behind at every size, 0.14x to 0.41x, for a structural reason: DataFusion's TopK keeps the best 100 rows of each partition while its input streams, so at 50,000,000 rows the whole query takes 6.21 ms, and MetalExec spends 7.72 plus 9.83 ms collecting and importing before the GPU does anything. Filters are behind for the same reason, and both are off by default; ArrowMetalConfig::all() switches them on.

plan nodewhat ArrowMetalConfig::default() does, and why
Full sort, ORDER BY without LIMITTaken from 250,000 exact rows; 6.9x to 28.8x ahead.
Top-k, ORDER BY … LIMITLeft, 0.14x to 0.41x: DataFusion's TopK keeps 100 rows per partition as it streams; MetalExec collects every row first.
GROUP BY, DISTINCTTaken for ten measured shapes over a MemTable of 50,000,000 rows or more; the rest left.
WHERELeft, 0.31x to 0.79x in memory, 0.51x to 0.55x over Parquet.
Hash join, inner, left, rightTranslated; the measured join table takes none of its 18 cells, as warm wins of up to 2.3x do not survive an idle GPU.
A sort above a filter, join or aggregateLeft: its row count is an estimate; the threshold needs an exact one.
From DATAFUSION.md, "What the default takes" and "What it leaves, and why". Every left node is a LEFT line in the report with this reason.

The CPU column is the part I would ask about first. DataFusion sorts 50,000,000 rows on 16 partitions in 1,581 ms of wall time and 6,104 ms of CPU time; with the rule the same query is 69.74 ms and 126 CPU-ms, because the GPU's time is not CPU time and unified memory never copies the columns across a bus. Of the 69.74 ms, 7.31 wait for the 6,104 input batches, 10.21 write them into Metal buffers, 48.21 are the sort.

A group-by is decided twice

The rule can see an aggregate's row count but not its number of groups, and the same count(*) is ahead at 1,000,000 groups and behind at 200. So an aggregate is decided twice. At plan time the rule describes the shape, the aggregate family, the keys and their type class, whether it reads a MemTable directly, and replaces the node only when the measured table takes that shape at some group count at the input's row count. When the query runs, MetalExec reads the first 262,144 rows, estimates the number of groups from a stratified sample of the keys with the bias-corrected Chao1 estimator, and looks the shape up again. If the table takes it, the GPU runs it; otherwise the batches already read and the rest of each partition stream into DataFusion's own aggregates, and the report says HANDBACK.

Which shapes the table takes is the strictest part of the crate. Warm, 38 of the 120 measured series are at least 1.65x ahead. But a GPU that has idled runs its next job two to four times slower than warm, and a query engine is idle most of the time, so each series must also win idle against idle, after 500 ms and after 5 s. Ten remain, all at 50,000,000 rows.

rulewhat it requires, and what survived it
Ahead warm, by enoughAt least 1.65x faster than DataFusion alone in its worst case, at two sizes or more: 38 of 120 series.
Ahead after the GPU idlesIts first run after 500 ms of idle, and after 5 s, at least as fast as DataFusion's first run after the same idle: removes 23.
A cheap hand-backWhen the run-time count says no, the first batches and the estimate cost at most 3%: removes 5.
What is leftTen series, count(*) and DISTINCT over two int32 keys and integer MIN/MAX over one int64 or two integer keys, from 50,000,000 rows: warm 2.31x to 4.27x; idle against idle 1.24x to 2.17x after 500 ms, 1.11x to 1.68x after 5 s. The hand-back measures 0.98x to 1.03x.
From DATAFUSION.md, "Aggregates". The warm and idle figures are the 20 cases the default runs on the GPU, Apple M4 Max, 50,000,000 rows, both table layouts, refit_check 2026-10-02 and default_import 2026-10-01.

DataFusion's answer, to the bit

A sort on a GPU is easy to make fast and easy to make subtly wrong. DataFusion orders floats in IEEE 754 totalOrder, a NaN with its sign bit set below negative infinity and -0.0 before +0.0, and compares strings by bytes. The GPU sort takes both as per-key options, so DataFusion's order comes out of the one radix pass. The differential grid is what made me trust it: 460 queries over 24 table configurations, 0 to 20,000 rows, null fractions of 0, 0.1 and 1.0, one partition and three, each run with and without the rule and compared bit for bit.

where a GPU could answer differentlywhat the rule does so the answer is DataFusion's
Float keysIEEE 754 totalOrder, as arrow-rs sorts: -NaN below -inf, -0.0 before +0.0, +NaN above +inf; descending the mirror.
NullsNULLS FIRST or LAST as written, inside the GPU sort; no extra key or pass.
StringsByte order; Utf8, LargeUtf8 and the Utf8View DataFusion reads from Parquet.
A float group key with several NaN bit patternsHanded back: DataFusion keeps each pattern as a group; the GPU would merge them.
Grouped float MIN/MAX with a NaN, or both zeros and a zero resultHanded back: DataFusion's answer depends on the order its rows arrive in.
From DATAFUSION.md, "DataFusion's semantics". The grid's hand-backs on data: 128 float MIN over a group with a NaN, 40 keys with several NaN patterns.

Run on 1 October: 13,632 query pairs, 11,931 with a node replaced, 0 mismatches, 0 hand-backs on an ArrowMetal error, 168 on data only DataFusion answers exactly. The grid also records what DataFusion does on strange data: its grouped float MIN starts each group at the largest finite value, so a group of positive infinities has a minimum of f64::MAX, and the rule reproduces that. Hand-back is the safety net too: on an ArrowMetal error, or when the memory pool refuses the collected input, MetalExec runs the subtree it replaced and returns DataFusion's result with a FALLBACK line. Nothing in the grid or the benchmark took that path.

Limits, plainly

  • MetalExec collects its input before it runs and has no spill path of its own; the batches count in DataFusion's memory pool, and a refusal hands the node back.
  • One GPU job at a time per process; the first GPU query of a process compiles its pipelines, 42 to 66 ms at 10,000,000 rows.
  • Aggregates only over a MemTable scan of 50,000,000 rows or more; joins, none by default.
  • DataFusion 55.1.0 exactly, arrow-rs 59.3.0, macOS on Apple silicon. Every number is from one M4 Max; the crate's examples/bench.rs measures yours.

The rest of 0.4.0

The DataFusion crate is the headline, not the whole release. The Polars MetalEngine now passes each sort key with Polars' null placement and float order, so sort().head() and top_k run as the GPU top-k: the top 100 by a nullable Float64 key over 50,000,000 rows went from 129.53 to 11.07 ms. Its group-by crossovers were refitted; the default takes 75 of 220 benchmark pairs, each 1.52x to 11.54x faster than the faster Polars engine, none behind. Every query the DuckDB rewrite's auto mode rewrites is 1.09x to 5.03x faster than DuckDB. The CHANGELOG has the rest, each with its file.

ArrowMetal is an independent Apache-2.0 project that implements Apache Arrow; Apache Arrow and Apache DataFusion are trademarks of the Apache Software Foundation. datafusion-arrowmetal = "0.4.1" beside datafusion = "=55.1.0", register the rule, read the report.

FAQ: DataFusion, the GPU and Apple silicon

Can Apache DataFusion use a GPU?

Not on its own, but it runs any physical optimizer rule you register. ArrowMetal 0.4.0's datafusion-arrowmetal crate is one for the Apple silicon GPU.

Does DataFusion run on the Apple silicon GPU with Metal?

With the rule, its full sorts and ten measured group-by shapes do; every other node stays DataFusion's. macOS on Apple silicon only.

Does the rule change DataFusion's results?

No. 13,632 query pairs run with and without it, floats compared bit for bit: 0 mismatches. Data only DataFusion answers exactly is handed back.

How much faster is DataFusion with the GPU?

On an Apple M4 Max, full sorts of 250,000 to 50 million rows are 6.9x to 28.8x faster than DataFusion alone. Top-k and filters are behind and left to DataFusion.

References & Citations

Subscribe to new posts from theaivibe.org

No spam — just new posts. One-click unsubscribe.
Share this article

Related Posts

sqljev: TypeSafe Jev's jev() for SQL Server, Postgres, Snowflake, BigQuery and DuckDB, on Jev or Open-Weight Laya
Data Engineering10 min read

sqljev: TypeSafe Jev's jev() for SQL Server, Postgres, Snowflake, BigQuery and DuckDB, on Jev or Open-Weight Laya

SQL cannot say 'the customer threatens to cancel'. sqljev adds jev(), jev_prob() and jev_choice() to SQL Server, PostgreSQL, MySQL, Snowflake, Databricks, BigQuery, Redshift and DuckDB, and answers them with Laya, an open-weight decision model that runs on your hardware, returns calibrated probabilities instead of text, and can be fine-tuned on your own tables. 140,000 decisions in 271 seconds on one laptop GPU; re-running all 13 queries, 0.9 seconds. Apache 2.0, version 0.1.0.

81 views
Read
pankhllm: The LLM Gateway That Learns to Skip the LLM, Without Replacing the Stack You Already Run
Data Engineering10 min read

pankhllm: The LLM Gateway That Learns to Skip the LLM, Without Replacing the Stack You Already Run

Most of an agent's LLM calls are not writing anything. They are decisions: which tool, which skill, which parameters, made thousands of times a day by a model paid in seconds and tokens. pankhllm sits where your app already calls an LLM, learns those decisions from its own traffic, and starts making them in 0.2 ms on a CPU with a 262 KB model. What it is not sure about still goes to your LLM. On the same 14 questions: a 12B planner 1,743 ms, Laya 49 ms, pankhllm's own model 4 ms, all 14 correct. Here is what it is, what it is not, and where it stops.

91 views
Read
Your Apple Silicon GPU Loses to One CPU Core Until a Million Rows. I Measured 111 Operations, Then Rebuilt ArrowMetal 0.2.0 Around the Answer
Data Engineering19 min read

Your Apple Silicon GPU Loses to One CPU Core Until a Million Rows. I Measured 111 Operations, Then Rebuilt ArrowMetal 0.2.0 Around the Answer

Every GPU data library benchmarks itself at 50 million rows. Your dataframe has 80,000. On an Apple M4 Max, summing 1,000 integers takes the GPU 112 microseconds and Polars less than one: the GPU is more than 100 times behind. I built one of these libraries, so I measured the row count where the GPU overtakes the fastest CPU code for 111 operations: the median needs 10,000,000 rows against a multi-core library, about a million against one core, and sixteen never get there. So ArrowMetal 0.2.0 refuses the GPU below the line, byte-identical, and around that router it grew GPU readers for CSV, JSON, nested Parquet, Delta Lake and Iceberg, a Polars engine and a DuckDB optimizer extension.

89 views
Read