Dask

Dask: Distributed Python at Scale

Dask (github.com/dask/dask (https://github.com/dask/dask 13,931 ), BSD-3-Clause, 2026.8.0 on 24 Aug 2026, pip 21,050 install "dask[dataframe]") parallelizes pandas 16,086 , NumPy and plain Python by splitting work into a graph of tasks. A Dask DataFrame is many pandas DataFrames; the same code runs on local threads or, with dask.distributed, on a cluster of worker processes, and Dask now optimizes DataFrame queries before running them.

bench/bench_dask.py: pandas-style code, partitionedPython
"""Dask DataFrame (local threads): the BookNest benchmark. Usage: bench_dask.py T1|T2|T3"""
import sys, time
import dask
import dask.bag as db
import dask.dataframe as dd
import json
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 = dd.read_parquet(f"{D}/{LINES}.parquet/*.parquet",
                            columns=["order_ts", "status", "book_id", "customer_id", "qty",
                                     "unit_price"])
    lines = lines[lines["status"] == "delivered"]
    books = dd.read_parquet(f"{D}/books.parquet/*.parquet", columns=["id", "genre"])
    customers = dd.read_parquet(f"{D}/customers.parquet/*.parquet",
                                columns=["customer_id", "country"])
    df = lines.merge(books.rename(columns={"id": "book_id"}).astype({"book_id": "int32"}),
                     on="book_id").merge(customers, on="customer_id")
    df["revenue"] = df["qty"] * df["unit_price"].astype("float64")
    df["month"] = df["order_ts"].dt.to_period("M").dt.to_timestamp()
    out = (df.groupby(["month", "genre", "country"])
           .agg(revenue=("revenue", "sum"), lines=("revenue", "count")).reset_index())
else:
    def explode(o):
        return [(o["order_ts"][:7], i["qty"] * i["unit_price"]) for i in o["items"]]
    bag = (db.read_text(f"{D}/raw/orders.jsonl", blocksize="64MiB").map(json.loads)
           .filter(lambda o: o["status"] == "delivered").map(explode).flatten())
    out = (bag.to_dataframe(meta={"month": "str", "revenue": "f8"})
           .groupby("month")["revenue"].sum().reset_index())
with dask.config.set(scheduler=sys.argv[2] if len(sys.argv) > 2 else "threads"):
    res = out.compute()
res.to_parquet(f"/home/dev/v7-l3/bench/out/dask_{sys.argv[1]}.parquet")
print(f"dask {sys.argv[1]} rows={len(res)} revenue={res['revenue'].sum():.2f} "
      f"query={time.perf_counter() - t0:.2f}s")

The JSON task uses a Dask bag of Python dictionaries, the flexible but slow path, because nested JSON has no fast columnar reader in pandas. Dask was 11 to 18 times slower than DuckDB 61,228 here because each partition still runs pandas, single-threaded per task, with Python objects for strings. Its strength is elsewhere: scaling existing pandas and NumPy code, arbitrary Python task graphs, and clusters (Kubernetes 5,150 , YARN, HPC schedulers, or the commercial Coiled) without a JVM.