| 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 |

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