Skip to content

benchmarks

Randy Zwitch edited this page Sep 22, 2026 · 1 revision

CPU benchmarks

pixi run bench runs benchmarks/bench_suite.mojo and prints CSV: one line per workload with rows, null density, kernel width, group shape, batch size, repetitions, best and mean nanoseconds, rows per second of the best run, and a checksum. Inputs are generated from a fixed seed before timing; each workload is evaluated once untimed, then timed five times, and results are validated outside the timed region (scalar and SIMD arithmetic must produce identical checksums).

  • BENCH_SMOKE=1 (pixi run bench-smoke) uses tiny inputs and one repetition. CI runs it as a correctness check; no timing thresholds are enforced on shared runners.
  • BENCH_LARGE=1 adds 10,000,000-row inputs.
  • Peak memory is measured externally: /usr/bin/time -v pixi run bench on Linux, /usr/bin/time -l on macOS. The suite does not count allocations; that needs an allocator hook Mojo does not expose, and timing is not used as a proxy.
  • The header line records the configuration, including the worker-thread limit (DATAFRAME_THREADS, default: physical cores). Run with DATAFRAME_THREADS=1 for single-threaded numbers.

Other focused benchmarks: bench-csv, bench-concat, bench-sort, bench-group-by, and bench-join.

Baseline (Mojo 1.1.0, AMD Ryzen Threadripper 3970X, Linux x86-64)

mojo=1.1.0 seed=20260918 physical_cores=32 native_f64_simd_width=4 threads=1 batch_size=1024

workload rows nulls kernel groups best ms rows/s (M)
arithmetic_chain 1,000 none scalar 0.06 16.3
arithmetic_chain 1,000 none simd4 0.06 17.4
nullable_compare 1,000 none simd4 0.03 37.7
filter 1,000 none simd4 0.03 28.7
global_sum 1,000 none scalar 0.01 80.5
global_count 1,000 none scalar 0.01 105.2
arithmetic_chain 1,000 10pct scalar 0.06 16.8
arithmetic_chain 1,000 10pct simd4 0.06 18.0
nullable_compare 1,000 10pct simd4 0.03 38.5
filter 1,000 10pct simd4 0.03 30.2
global_sum 1,000 10pct scalar 0.01 82.9
global_count 1,000 10pct scalar 0.01 105.9
grouped_sum_count 1,000 10pct scalar low 0.04 22.4
grouped_sum_count 1,000 10pct scalar high 0.07 14.7
grouped_sum_count 1,000 10pct scalar skewed 0.06 15.9
arithmetic_chain 100,000 none scalar 6.92 14.4
arithmetic_chain 100,000 none simd4 6.50 15.4
nullable_compare 100,000 none simd4 2.69 37.2
filter 100,000 none simd4 3.60 27.8
global_sum 100,000 none scalar 1.17 85.4
global_count 100,000 none scalar 0.88 113.6
arithmetic_chain 100,000 10pct scalar 6.71 14.9
arithmetic_chain 100,000 10pct simd4 6.30 15.9
nullable_compare 100,000 10pct simd4 2.61 38.3
filter 100,000 10pct simd4 3.38 29.6
global_sum 100,000 10pct scalar 1.11 90.4
global_count 100,000 10pct scalar 0.85 117.7
grouped_sum_count 100,000 10pct scalar low 3.79 26.4
grouped_sum_count 100,000 10pct scalar high 6.98 14.3
grouped_sum_count 100,000 10pct scalar skewed 6.42 15.6
arithmetic_chain 1,000,000 none scalar 183.28 5.5
arithmetic_chain 1,000,000 none simd4 191.88 5.2
nullable_compare 1,000,000 none simd4 37.14 26.9
filter 1,000,000 none simd4 46.02 21.7
global_sum 1,000,000 none scalar 11.66 85.8
global_count 1,000,000 none scalar 8.82 113.3
arithmetic_chain 1,000,000 10pct scalar 184.45 5.4
arithmetic_chain 1,000,000 10pct simd4 198.28 5.0
nullable_compare 1,000,000 10pct simd4 35.85 27.9
filter 1,000,000 10pct simd4 44.09 22.7
global_sum 1,000,000 10pct scalar 11.48 87.1
global_count 1,000,000 10pct scalar 8.99 111.2
grouped_sum_count 1,000,000 10pct scalar low 42.38 23.6
grouped_sum_count 1,000,000 10pct scalar high 55.54 18.0
grouped_sum_count 1,000,000 10pct scalar skewed 53.84 18.6

After fixing quadratic batch reassembly

The baseline exposed a bug: appending each evaluated batch to the output reserved exactly the new length, so every append reallocated and copied the whole column, making reassembly quadratic in the number of batches. With geometric growth, the 1,000,000-row workloads became:

workload nulls before ms after ms
arithmetic_chain (simd4) none 191.9 58.3
arithmetic_chain (scalar) none 183.3 61.8
nullable_compare none 37.1 26.7
filter none 46.0 36.1
global_sum / global_count none 11.7 / 8.8 11.8 / 9.0

The analysis below still holds: per-node materialization and copying dominate the remaining arithmetic time, and SIMD is only slightly faster than scalar.

After elementwise fusion (#4)

Fused Float64 subtrees evaluate in one SIMD pass reading source buffers directly. At 1,000,000 rows the 4-lane arithmetic chain fell from 58.3 ms to 22.3 ms (183 ms at the original baseline), and the nullable comparison from 26.7 ms to 18.8 ms; checksums are unchanged. The 1-lane path gains nothing, since per-row program interpretation replaces per-node materialization.

After shared buffers (#34)

Columns became windows onto reference-counted buffers, so per-batch slices and projections no longer copy. At 1,000,000 rows (no nulls):

workload before ms after ms
global_sum 11.8 7.2
global_count 9.0 4.5
grouped_sum_count (low cardinality) 43 28
filter 36.1 38.7
arithmetic_chain (simd4) 22.3 24.0
nullable_compare 18.8 20.5

Scans that only read slices gain the most. Fused arithmetic, comparisons, and filters pay roughly 5-8% for the extra indirection (offset plus shared-pointer dereference) on per-row validity reads; hoisting buffer references in those kernels, or word-at-a-time validity, is the planned follow-up.

After contiguous UTF-8 strings (#33)

String columns moved from List[String] to one UTF-8 buffer with Int64 offsets (Arrow large_utf8). Throughput on the string workloads is unchanged within noise; the per-row work (hashing, CSV tokenizing, rank merging) dominates, not string storage:

workload before ms after ms
CSV read, 100k rows with a quoted string field 92.7 94.7
group_by string + int key, 16 / 10k / 111k groups 10.7 / 12.5 / 27.2 10.6 / 12.6 / 28.6
sort one string key (ranked), 200k rows 90.0 91.5
sort one string key (comparator reference), 200k rows 120.9 81.8

The comparator sort gains because comparisons read borrowed slices instead of copying Strings. Memory falls: peak RSS for a 2M-row column plus a 1M-row take is 104 → 92 MB for ~10-byte keys (which String already stores inline) and 231 → 156 MB for ~30-byte strings, since each row costs its bytes plus an 8-byte offset instead of a 24-byte String and a separate heap allocation.

What the baseline shows

  • The five-node Float64 arithmetic chain runs at about 5 million rows/s, and the 4-lane SIMD kernel is no faster than the 1-lane one. Arithmetic is not the bottleneck: every node materializes a new batch Series, source columns are sliced by copying values and rebuilding validity one element at a time, and batch results are appended into the output. Borrowed column views (#3) and fused elementwise kernels (#4) target exactly this.
  • A single comparison runs at about 27 million rows/s and a filter at about 22 million, both again dominated by per-batch copies and take.
  • Global sums and counts reach 85-110 million rows/s: they read slices and update one state, with little materialization.
  • Grouped sum+count is 18-24 million rows/s and slows with cardinality, as expected from hashing and per-group state.
  • Null density (10%) barely changes any workload, because validity is handled per element either way.

No workload here claims a speedup; these numbers are the reference for the execution changes in #3-#8.

Parallel reductions (#6, #8)

Reductions split rows into one contiguous partition per worker thread; each worker reduces its partition into private state, and states merge in partition order. At 1,000,000 rows (32-core Threadripper; default workers = min(physical cores, rows / 65536) = 15):

workload 1 thread default
global_sum 5.2 ms 1.4 ms 3.7x
global_count 3.3 ms 1.7 ms 1.9x
grouped_sum_count, 16 groups 26.8 ms 22.0 ms hashing (serial) dominates
grouped_sum_count, ~100k groups 40.2 ms 38.2 ms workers capped by group count

Grouped reductions give each worker state for every group, so the worker count is capped at rows / (4 x groups). Hash grouping itself, row-wise expressions, and filters are still single-threaded (#5, #7).

Direct SIMD loads in unfused float kernels (#3)

Unfused float kernels (Float32, %, //, **, clip, unary math, and any expression fusion does not cover) now load contiguous vectors straight from the shared column buffers instead of gathering lane by lane, and apply validity as a vector mask. Single-threaded, 1,000,000 rows:

expression before after
Float64 x % 7 20.0 ms 13.8 ms
Float32 y * 2 + 1 32.6 ms 22.0 ms
Float64 sqrt(x) 15.5 ms 10.6 ms

Remaining allocations per batch are the output values, the validity list, and one intermediate per unfused node; source windows allocate nothing.

Row-parallel expressions and filters (#5, #7)

Row-shaped expressions evaluate one contiguous, batch-aligned partition per worker and join the pieces in order. Filters compact in parallel: workers collect true-row indices per partition, then gather rows into disjoint output ranges aligned to 8 rows (no shared validity bytes). 1,000,000 rows, no nulls:

workload 1 thread default (15 workers)
arithmetic_chain (fused SIMD) 23.5 ms 6.3 ms
arithmetic_chain (scalar) 63.7 ms 15.0 ms
nullable_compare 19.9 ms 4.7 ms
filter (end to end) 37.8-42.4 ms 10.7 ms

Crossover: below 2 x 65,536 rows everything stays single-threaded, which avoids thread startup (~20-50 us per worker) on small frames. Joining the per-worker pieces is a serial O(n) copy and is the main remaining cost at high worker counts.

Bit-packed Booleans (#80)

Boolean values moved from one byte per row to one bit, so Boolean results (comparisons, predicates, filter masks) write 8x fewer bytes. 1,000,000 rows, default workers:

workload before after
nullable_compare 4.7 ms 2.7 ms
filter (end to end) 10.7 ms 8.7 ms

Peak RSS for five 20,000,000-row Boolean columns: 252 MB -> 204 MB (the remainder is the transient List[Bool] inputs).

Reading a column out of the library (#91)

Summing a 1,000,000-row Float64 column with 10% nulls, compiled with mojo build. All three paths produce the same total:

how the consumer reads time
Series.get(i) -> AnyValue per element ~3,500 us
Column.value(i) (bounds check, raises) ~4,740 us
unsafe_values() + unsafe_validity() ~700 us

AnyValue carries a String field a numeric column never fills, and Series.get rescans the dtype list per call; Column.value is slower still because it bounds-checks and re-reads the validity bit on every element. The buffer path resolves the dtype once and then walks the payload.

Reading through the buffers is the intended way for another package to consume a column. Series.numeric[D](), .string() and .bool() already hand back a window that shares the buffers in O(1), so the cost is one pointer dereference per row, not a copy.

Measure with an optimized build: under mojo run the difference largely disappears.

Hash-partitioned grouping (#104)

Grouping encoded keys serially into dense ids (17-22 ms of a 22-57 ms group-by at 1M rows), then reduced in parallel with every worker holding state for every group, so the worker count collapsed as cardinality rose. Two attempts to parallelize the encoding by merging worker dictionaries (#8) lost at high cardinality, because that merge costs workers x distinct keys.

Rows are now hashed by key and stably partitioned into buckets, the key and aggregated columns are gathered into bucket order with take_parallel, and each bucket runs the existing encode-and-reduce on its own. Equal keys always share a bucket, so nothing is merged. maintain_order=True costs one O(groups) sort by first input row afterwards.

Partitioning is not free: the gather is a full materialization of the frame, and it only repays when the serial encode it replaces is expensive. low_cardinality decides beforehand by hashing a 4,096-row strided sample on the calling thread and counting how much of hash space it touches, which costs far less than the pass it guards.

1,000,000 rows, bench-polars, 32 threads, best of 5, against main on the same machine back to back:

workload main partitioned vs Polars, before -> after
grouped_high (~100k groups) 51.72 ms 24.47 ms 4.2x -> 2.0x
grouped_low (16 groups) 23.21 ms 23.63 ms 1.8x -> 1.8x
grouped_skew 23.90 ms 25.00 ms 1.6x -> 1.6x
grouped_str (100 groups) 34.99 ms 33.85 ms 2.3x -> 2.2x

High cardinality is 2.1x faster, which is where the earlier attempts lost, and the gate keeps the low-cardinality shapes at parity.

At 200,000 rows (bench_group_by, best of 3, three rounds) the picture is the same but smaller, because the fixed costs are a larger share:

keys groups main ms partitioned ms
1 16 5.0-5.4 5.2-5.3
3 16 14.3-14.6 14.6-15.0
1 10,000 7.2-7.6 6.1-7.0
2 10,000 12.9-13.6 14.3-15.3
1 110,609 17.2-17.4 12.3-13.0
2 110,609 28.8-28.9 20.9-22.7
3 110,609 37.4-37.8 23.1-23.9

One shape regresses: 10,000 groups on two keys, about 10%, where the sample says partition but the gather is not repaid at that cardinality and frame size. A reusable thread pool would likely remove it, since the extra hash, scatter and gather stages each pay thread creation today; that was measured and could not ship (#103).

Parallel joint-key encoding for joins (#105)

An inner join of 1,000,000 rows to 500,000 distinct keys took 235 ms, of which encode_rows over both sides stacked (1.5M rows) was 204 ms. The dictionary holds every distinct key of both sides at once, so at high cardinality it spends most of its time missing cache -- the same problem #104 fixed for grouping.

A join's row order comes from iterating rows, not from the id numbering: ids exist only so equal keys compare equal. So the ids may be assigned in any consistent order, and encode_partitioned assigns them one hash bucket at a time, each bucket with its own small dictionary and an offset. The join's semantics are untouched.

Separately, the rows belonging to each key id were held as one List per distinct key, which is a heap allocation per key. They are now a flat CSR layout built with a counting sort, which preserves the increasing-row-order match sequence the contract documents.

1,000,000 x 500,000 inner join, best of 3:

ms
before 235
partitioned encode 150
plus flat group index 134

In bench-polars (more output columns, so assembly weighs more) the same join goes 283 -> 188 ms, and from 10.5x to 7.0x of Polars.

What is left is the part this did not touch: the match loop and the output gather are still serial, and each bucket pair is not joined independently. That is the rest of #105.

Parallel stable sort (#108)

sort, arg_sort, top_k and bottom_k ran a single bottom-up mergesort over row indices. Sorting now splits rows into one contiguous range per worker, sorts each range, and merges the runs in rounds.

Stability is preserved by construction rather than by luck: ranges are cut in row order, each range sort is the same stable mergesort as before, and every merge prefers the earlier run when the comparator reports a tie. The result is therefore the serial order for any worker count, which the tests assert row by row rather than by checking sortedness.

1,000,000 rows, two keys, bench-polars, best of 5:

ms vs Polars
before 594 11.9x
ranges sorted in parallel 446 7.6x
plus the fixes below 346 6.2x

Two costs were worth more than the parallelism itself:

  • Each merge round copied the whole index array to hand workers a source buffer; they now read it by address and write disjoint output ranges.
  • _rank_less walked a List[List[Int]] per comparison, and a merge calls it once per output row. Sorting by one key, the common case, now compares that key directly.

What remains is the merge tree's shape: with 15 runs the last round is a single job merging the whole array, so the tail of the sort is serial. A parallel multiway merge, where each worker binary-searches splitters to claim an output slice, would remove it.

Float parsing in CSV (#107)

Reading 1,000,000 rows of 8 columns took 1,260 ms against Polars' 15 ms. Three attempts to speed this up by reasoning about the code failed, all measuring within noise of main: removing the per-field String, a SIMD structural scan that bulk-copies runs of ordinary bytes, and replacing the per-field Variant dispatch with an integer tag. Guessing does not work here.

Ablation does. Replacing one parser at a time with a constant, on the 4-column bench_csv file:

variant best ms
unchanged 79.6
integer parsing skipped as well 50.5
float parsing skipped 51.9

So float parsing alone is about 28 ms of 79 ms, and integer parsing about 1.3 ms. One of the four columns is a Float64, so that is roughly 140 ns per float field against 6.5 ns per integer field.

The reason is in parse_float64: Mojo's own parser accepts things CSV must reject (it reads 2024-02-28 as a number), so the grammar was checked with a full scan and the text was then converted with a second full scan.

Plain decimals now take one pass that validates and computes together, and the result is exact rather than approximate: a mantissa below 253 divided by a power of ten below 1023 is two exact operands, so the division is correctly rounded once. Anything else -- a sign, an exponent, nan, an infinity spelling, more than 19 digits, more than 22 fraction digits, a mantissa too large to hold exactly -- falls through to the original parser, so the accepted grammar and every value are unchanged.

100,000 rows, 4 columns best ms
before 80.5-80.9
after 65.8-67.5

End to end on the 8-column file: 1,260 -> 1,081 ms.

This is 1.2x, not the order of magnitude the workload needs. What it buys is a correct reading of where the time goes: per field, roughly 36 ns was parsing and 31 ns is the surrounding plumbing, against about 4 ns per byte of input. Parsing is now much cheaper, so the plumbing and the single-threadedness are what remain.

Parallel CSV decoding (#107)

Reading 1,000,000 rows of 8 columns took 1,081 ms on one core. Blocks are now split at record boundaries and decoded on worker threads.

best ms vs Polars
before 1,081 86.7x
parallel ranges 215 13.4x
plus a SIMD boundary scan 161 11.1x
plus bulk column appends 145 11.5x

Timed by phase on the 1M-row, 50 MB file, so the remaining work is not a guess:

phase ms
file read 19
boundary scan 7
parallel decode 107-131
concatenating 64 partial frames 51

The boundary scan was about 200 ms before it was summarised by block, and is now negligible. Concatenation was 69 ms until Column._append_column stopped appending one bounds-checked element at a time and copied the window in bulk; that is on the path of concat, vstack, batch reassembly and every parallel stage's reassembly, not only this reader. What remains of it is the string column's bytes and offsets, and the bytewise validity merge.

Decoding is now the bulk, it is already spread across workers, and what is left inside it is the ~31 ns of per-field work that the ablation in the float section measured.

Once decoding was parallel, the boundary scan was the serial remainder: it ran a byte at a time over every block, at roughly the 4 ns/byte the tokenizer costs, which is most of 215 ms for a 50 MB file. A block that contains no quote cannot change parity, so every newline in it is a record boundary and the block can be summarised -- one SIMD load, a quote test and a newline count -- instead of walked. Only blocks holding a quote, a split target, or the final boundary need the byte loop, and that loop remains the definition of the format.

The 100,000-row, 4-column bench_csv file goes 66 -> 18 ms, so small files gain too rather than paying for the machinery.

Why splitting is the whole problem

A block boundary can land inside a quoted field, and a quoted field may contain newlines, so a worker cannot scan forward to the next newline and call it a record start. record_splits decides by quote parity: a newline is a boundary only when an even number of quotes precede it. Toggling on every quote byte handles CSV's doubled-quote escape with no special case, because "" toggles twice.

Blocks are read sequentially and bytes after the last complete record are carried into the next block, so memory stays bounded by the block size (8 MiB for parallel reads) rather than by the file.

Options that depend on counting records from the start of the file -- n_rows, skip_rows, comment_prefix, ignore_errors, truncate_ragged_lines -- keep the serial path, because a range cannot resolve them on its own.

What the tests had to catch

The first version sized workers with worker_count, which measures rows, against a byte count. A 64 KiB block yielded one worker, so the parallel path never ran and every CSV test passed while testing nothing. The differential tests caught it by asserting the parallel path was reachable before comparing. Worker sizing is by bytes now, a quarter megabyte each.

Two further defects surfaced only once the path was live, both from the existing contract suites: an empty quote character indexed byte 0 of an empty string and crashed, and every range looked for a byte-order mark, which is meaningful only at the very start of a file.

Clone this wiki locally