Skip to content
woozykingPublic

About

Datapunk is the "Cyberpunk 2077" of data analytics: a benchmark suite for data engines with real-world workloads

Resources

Stars

0 stars

Watchers

0 watching

Forks

Repository files navigation

Datapunk

Synthetic benchmarks like TPC-H offer idealized metrics. Just as PC gamers reach for Cyberpunk 2077 to stress a rig, data engineers need a real-world suite to see how engines behave on practical analytics — and, crucially, on a small box.

Datapunk is that benchmark "game."

Engines

  • Pandas - baseline, no memory configuration
  • Polars — default in-memory + streaming engine twin
  • DuckDB — default + spill-to-disk twin (explicit memory_limit)
  • Dask — threaded scheduler
  • Daft — streaming architecture, no memory configuration

Each engine runs with out-of-the-box defaults. Where the engine supports memory configuration, a "twin" variant is included that receives an explicit memory budget and spill directory — same query, with and without the engine knowing its constraint.

The two-run model

Every suite runs each engine twice:

  1. Small, uncapped — a single month (~3M rows). No memory limit. This is the apples-to-apples comparison of speed and correctness when everything fits in RAM.
  2. Large, capped — the full 24-month window (30M+ rows) under a physical memory ceiling. This surfaces out-of-core behaviour: engines that stream and spill survive; engines that must materialise everything OOM — and the OOM is recorded as a first-class result, not a crash.

How the memory cap works

Each engine runs in its own isolated subprocess (fresh interpreter, no shared allocator state, no cross-engine import pollution). The parent enforces the cap with a physical-RSS watchdog: it polls the worker's resident memory and SIGKILLs it the instant it crosses the limit, recording an OOM.

Metrics

  • Velocity — speedup vs. the slowest successful engine in that run (best ★).
  • Peak RAM — peak resident memory (watchdog max ⨄ kernel high-water mark).
  • Stress — peak RAM ÷ cap (large run only; ≥1.0 / OOM means it hit the ceiling).

Verification compares an order-independent fingerprint (row count + per-column aggregate stats) across engines with floating-point tolerance — robust to legitimate tie-break and last-ULP differences that a strict row-by-row check flaked on.

Layout

datapunk/         # installable package
  config.py       # dataset source, default cap, run modes
  dataio.py       # download, schema unification, page-cache warming
  runner.py       # subprocess + RSS watchdog (the cap)
  worker.py       # isolated per-engine entry point
  fingerprint.py  # tolerant cross-engine consistency check (NOT DuckPipe's
                  # own fingerprint.py -- same name, unrelated job: this one
                  # verifies engines agree, DuckPipe's decides what to re-run)
  reporter.py     # two-run orchestration, scorecard, JSON export
notebooks/        # one suite per notebook: pipeline-shape diagrams, charts, scorecard
suites/           # one dir per suite: each engine's real logic as a headless
                  # DuckPipe pipeline (see "Per-engine pipelines" below)
pipeline.py             # DuckPipe DAG: one cached task per notebook (see below)
show_pipeline_mermaid.py # renders the DAG below, each suite's real Pandas shape nested in
run_benchmark.py        # thin duckpipe.run(pipeline.py) entrypoint

Pipeline DAG, each suite's own real Pandas pipeline nested as a subgraph (uv run python show_pipeline_mermaid.py — plain duckpipe show pipeline.py --mermaid can't show this part; see "Orchestration" below for why it's Pandas specifically, and why the benchmark's own timing is never routed through this):

flowchart TD
    subgraph t_suite_01 ["suite_01"]
        t_suite_01__t_extract["extract"]
        t_suite_01__t_derive["derive"]
        t_suite_01__t_aggregate["aggregate"]
        t_suite_01__t_finalize["finalize"]
        t_suite_01__t_extract --> t_suite_01__t_derive
        t_suite_01__t_derive --> t_suite_01__t_aggregate
        t_suite_01__t_aggregate --> t_suite_01__t_finalize
        class t_suite_01__t_extract success
        class t_suite_01__t_derive success
        class t_suite_01__t_aggregate success
        class t_suite_01__t_finalize success
    end
    subgraph t_suite_02 ["suite_02"]
        t_suite_02__t_extract["extract"]
        t_suite_02__t_clean["clean"]
        t_suite_02__t_aggregate["aggregate"]
        t_suite_02__t_extract --> t_suite_02__t_clean
        t_suite_02__t_clean --> t_suite_02__t_aggregate
        class t_suite_02__t_extract success
        class t_suite_02__t_clean success
        class t_suite_02__t_aggregate success
    end
    subgraph t_suite_03 ["suite_03"]
        t_suite_03__t_extract_lookup["extract_lookup"]
        t_suite_03__t_extract_trips["extract_trips"]
        t_suite_03__t_join_aggregate["join_aggregate"]
        t_suite_03__t_extract_trips --> t_suite_03__t_join_aggregate
        t_suite_03__t_extract_lookup --> t_suite_03__t_join_aggregate
        class t_suite_03__t_extract_lookup success
        class t_suite_03__t_extract_trips success
        class t_suite_03__t_join_aggregate success
    end
    subgraph t_suite_04 ["suite_04"]
        t_suite_04__t_extract["extract"]
        t_suite_04__t_derive["derive"]
        t_suite_04__t_aggregate["aggregate"]
        t_suite_04__t_extract --> t_suite_04__t_derive
        t_suite_04__t_derive --> t_suite_04__t_aggregate
        class t_suite_04__t_extract success
        class t_suite_04__t_derive success
        class t_suite_04__t_aggregate success
    end
    subgraph t_suite_05 ["suite_05"]
        t_suite_05__t_compute_bounds["compute_bounds"]
        t_suite_05__t_extract_candidates["extract_candidates"]
        t_suite_05__t_finalize["finalize"]
        t_suite_05__t_compute_bounds --> t_suite_05__t_extract_candidates
        t_suite_05__t_extract_candidates --> t_suite_05__t_finalize
        class t_suite_05__t_compute_bounds success
        class t_suite_05__t_extract_candidates success
        class t_suite_05__t_finalize success
    end
    t_validate_dashboard_json["validate_dashboard_json"]
    t_suite_01 --> t_validate_dashboard_json
    t_suite_02 --> t_validate_dashboard_json
    t_suite_03 --> t_validate_dashboard_json
    t_suite_04 --> t_validate_dashboard_json
    t_suite_05 --> t_validate_dashboard_json
    class t_suite_01 success
    class t_suite_02 success
    class t_suite_03 success
    class t_suite_04 success
    class t_suite_05 success
    class t_validate_dashboard_json success
    classDef success fill:#d4f7dc,stroke:#2f9e44,color:#1a1a1a
Loading

Three of the five suites nest as plain chains of different lengths; suite_03 is a genuine two-root join (fact and dimension tables, merged only where the join needs both); suite_05 is a two-phase reduction (local top-N per file, then re-ranked once over the small combined set). Same rendering, just pointed at whatever each suite's own pipeline actually looks like.

Quickstart

uv run python run_benchmark.py            # skips notebooks that haven't changed
uv run python run_benchmark.py --force    # re-runs every notebook regardless

Results are written to docs/benchmark_results.json, used by the dashboard at https://woozyking.github.io/datapunk/.

Orchestration

pipeline.py wraps each notebook as a DuckPipe task, fingerprinted on the notebook's own content — so a notebooks/03_*.ipynb edit re-runs only suite 03, not all five. A notebook's actual logic lives in its .ipynb JSON rather than in the wrapper's own Python source, so DuckPipe's extra_fingerprint (its documented escape hatch for exactly this — see its own docs) is keyed on a hash of the notebook file instead of the wrapper function.

Each suite task also does one real, cheap, once-per-invocation nested duckpipe.run() against that suite's own Pandas pipeline (Pandas specifically because it has the most-decomposed shape across all five — a genuine two-root join for suite 03, a two-phase reduction for suite 05, plain chains for the rest, see each pipeline's own docstring). This is what show_pipeline_mermaid.py renders as a real nested subgraph per suite, using DuckPipe's own recursive to_mermaid(..., subgraphs=...). It is deliberately never the benchmark's own timed measurement, which stays exactly analytics.py's direct-call bypass, unchanged: a real duckpipe.run() call costs roughly 50–60ms per invocation (measured directly), enough to swamp the fastest engines' own measured time rather than just add noise to it — precisely the reason that bypass exists in the first place (see "Per-engine pipelines" below). State for this nested run lives in its own suites/suite_NN/ pipeline_shape.duckdb, gitignored the same way duckpipe.db is.

Run history (per-suite status, timing) lives in duckpipe.db next to this file — inspect it directly, the same as any DuckPipe pipeline:

uv run duckpipe show pipeline.py    # DAG + cache status per suite
uv run duckpipe stats duckpipe.db   # timing history

duckpipe.db is gitignored, deliberately — it's a binary file that changes on every run (a fresh run_id, new timestamps) even when nothing meaningful happened, so a git diff on it is never reviewable, and two branches each committing their own locally-run copy would be an unresolvable binary merge conflict. The published, diffable record is docs/benchmark_results.json; duckpipe.db only needs to exist for your own machine's incremental-rerun convenience, which doesn't require sharing it. A fresh clone (or CI) starts with no prior state and re-runs every suite once, same as it always did before this DAG existed.

Per-engine pipelines

Each suite's actual read/clean/aggregate logic lives in suites/suite_NN/*_pipeline.py — one small, headless DuckPipe pipeline per engine, not inline in the notebook. Two consequences:

  • Independently runnable and inspectable. uv run duckpipe run suites/suite_02/pandas_pipeline.py exercises one engine standalone, with its own real caching and state history. uv run duckpipe show suites/suite_02/pandas_pipeline.py --mermaid shows exactly what it does — and every notebook now renders this for all its engines in a "Pipeline shape" cell, instead of asking you to read Python to find out.
  • Task granularity follows each engine's own idiomatic shape, not a forced template. Pandas/Dask/Polars/Daft decompose into real stages (extract → clean → aggregate, sometimes a genuine fan-in for a join, or a two-phase reduction for top-N); DuckDB is deliberately one task everywhere — splitting a single declarative query into artificial stages would force intermediate materialization and misrepresent the query-planner-fused performance it's actually being measured for.

suites/suite_NN/analytics.py adapts these pipelines to DatapunkReporter.run_all's calling convention — each wrapper calls its pipeline's tasks directly (bypassing DuckPipe's DAG and state file entirely), because a benchmark needs a fresh, full-cost measurement on every timed iteration, never a cached skip. Confirmed directly, not assumed: an interleaved, same-process A/B timing test against the original inline functions showed no measurable difference (well within run-to-run noise), and the isolated per-call mechanism overhead is ~77 nanoseconds — about 0.0002% of even DuckDB's fastest aggregation.

The one real cost worth knowing about: importing a headless pipeline module pulls in DuckPipe, which (as of DuckPipe ≥0.2.0) no longer eagerly imports duckdb just for that — deferred until something actually opens a state file, which these direct-call wrappers never do. On DuckPipe 0.1.x this cost ~24MB of RSS per engine subprocess, even for engines with nothing to do with DuckDB; confirmed fixed by measuring peak RSS before and after upgrading.

About

Datapunk is the "Cyberpunk 2077" of data analytics: a benchmark suite for data engines with real-world workloads

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages