Polars 268,908 (github.com/pola-rs/polars (https://github.com/pola-rs/polars 39,914 ), MIT, 1.44.2 on 9 Sep 2026, pip 21,050 install polars) is a DataFrame library written in Rust on the Apache Arrow 129 memory format. Its lazy API builds a query plan like Spark 129 's, optimizes it (projection and predicate pushdown into Parquet 129 , join reordering), then executes it in parallel on all cores, with a streaming engine for inputs larger than memory. The benchmark's Polars version reads like PySpark 129 with fewer ceremonies:
"""Polars: the BookNest benchmark. Usage: bench_polars.py T1|T2|T3"""
import sys, time
import polars as pl
D = "/home/dev/v7-l3/ch05/data"
LINES = "order_lines_x10" if sys.argv[1] == "T3" else "order_lines" # T3: 10x rows
t0 = time.perf_counter()
if sys.argv[1] in ("T1", "T3"):
lines = pl.scan_parquet(f"{D}/{LINES}.parquet/*.parquet")
books = pl.scan_parquet(f"{D}/books.parquet/*.parquet").select(
pl.col("id").cast(pl.Int32).alias("book_id"), "genre")
customers = pl.scan_parquet(f"{D}/customers.parquet/*.parquet").select(
"customer_id", "country")
out = (lines.filter(pl.col("status") == "delivered")
.join(books, on="book_id").join(customers, on="customer_id")
.group_by(pl.col("order_ts").dt.truncate("1mo").alias("month"), "genre", "country")
.agg((pl.col("qty") * pl.col("unit_price").cast(pl.Float64)).sum().alias("revenue"),
pl.len().alias("lines"))
.sort("month", "genre", "country"))
else:
item = pl.Struct({"book_id": pl.Int32, "qty": pl.Int32, "unit_price": pl.Float64})
schema = {"order_id": pl.Int64, "order_ts": pl.String, "status": pl.String,
"items": pl.List(item)}
out = (pl.scan_ndjson(f"{D}/raw/orders.jsonl", schema=schema)
.filter(pl.col("status") == "delivered").explode("items").unnest("items")
.group_by(pl.col("order_ts").str.slice(0, 7).alias("month"))
.agg((pl.col("qty") * pl.col("unit_price")).sum().alias("revenue"))
.sort("month"))
df = out.collect()
df.write_parquet(f"/home/dev/v7-l3/bench/out/polars_{sys.argv[1]}.parquet")
print(f"polars {sys.argv[1]} rows={df.height} revenue={df['revenue'].sum():.2f} "
f"query={time.perf_counter() - t0:.2f}s")scan_parquet and scan_ndjson return lazy frames; nothing runs until collect(). Polars reads Spark's legacy INT96 timestamps as naive UTC values, so the months come out identical to Spark's with a UTC session. Its weak spots are the ones the benchmark exposed: JSON parsing is slower than DuckDB 61,228 's, and the 13.8-million-row join used about 1.9 GB of memory, twenty times DuckDB's footprint.