Benchmark Results

Benchmark Results and What They Mean

Median wall time of three fresh-process runs, ratio to DuckDB 61,228 , on the shared 4-CPU host
Engine T1 joins T2 JSON T3 at 10x Peak memory, T3
DuckDB 0.29 s (1.0x) 0.47 s (1.0x) 0.84 s (1.0x) 91 MiB
Polars 268,908 0.55 s (1.9x) 2.67 s (5.7x) 2.89 s (3.4x) 1,926 MiB
Daft 1.10 s (3.8x) 2.49 s (5.3x) 2.77 s (3.3x) 855 MiB
pandas 16,086 + PyArrow 129 2.47 s (8.5x) 6.80 s (14.5x) 24.59 s (29.3x) 4,829 MiB
Dask (threads) 3.22 s (11.1x) 8.58 s (18.3x) 14.85 s (17.7x) 2,079 MiB
Spark 129 local[4] 18.22 s (62.8x) 18.24 s (38.8x) 20.38 s (24.3x) 794 MiB
Median wall time per engine for T1 and T3, on one scale (fresh process per run)
Median wall time per engine for T1 and T3, on one scale (fresh process per run)

The single-node engines won by one to two orders of magnitude, and most of Spark's time was fixed cost: about 5 seconds to build a session and several more of JVM start-up, JIT warm-up and shutdown, nearly the same for 1.4 or 13.8 million rows. With the session already running, the picture changes:

bench/spark_warm.py: the T1 and T3 query three times in one Spark session (abridged)Python
for name in ("order_lines", "order_lines_x10"):
    for run in range(3):
        t0 = time.perf_counter()
        (spark.read.parquet(f"{D}/{name}.parquet").where("status = 'delivered'")
         .join(books, "book_id").join(customers, "customer_id")
         .groupBy(F.date_trunc("month", "order_ts").alias("month"), "genre", "country")
         .agg(F.sum(F.col("qty") * F.col("unit_price").cast("double")).alias("revenue"))
         .write.mode("overwrite").parquet(f"/home/dev/v7-l3/bench/out/warm_{name}"))
        print(f"{name:<16} run {run + 1}: {time.perf_counter() - t0:5.2f}s")
Output
order_lines      run 1:  6.41s
order_lines      run 2:  2.27s
order_lines      run 3:  1.33s
order_lines_x10  run 1:  3.94s
order_lines_x10  run 2:  3.86s
order_lines_x10  run 3:  3.10s

Warm, Spark ran T3 in about 3 seconds, on par with Polars and Daft from a cold start and only about 3.7 times DuckDB, and the gap keeps narrowing as data grows while the single-node engines approach the machine's memory and cores. The memory column matters as much as the times: pandas needed 4.8 GB for 13.8 million rows because it materializes everything as Python-friendly columns, and Polars 1.9 GB, while DuckDB streamed the join in under 100 MB. JSON (T2) separated the engines with native nested readers from those that fall back to Python dictionaries. These are small workloads on one shared machine: they show where the overheads lie, not how the engines compare on a cluster with terabytes.