Skip to content

docs(adr): ADR-004 out-of-core scan_parquet backend for polars /load - #1059

Open
paddymul wants to merge 1 commit into
mainfrom
adr-004-scan-parquet-backend
Open

paddymul wants to merge 1 commit into
mainfrom
adr-004-scan-parquet-backend

Conversation

@paddymul

@paddymul paddymul commented Oct 6, 2026

Copy link
Copy Markdown
Collaborator

Problem

The destination for polars /load is an out-of-core scan_parquet backend (#1000 and #1005, the eager stopgaps, are closed). No design for it is in the repo: ADR-003 lists "the scan_parquet backend (paging, search, sort, memory)" under "Not decided here", and the only implementation, #1011, is closed. On a 12M x 43 parquet file, anything that reads one or two columns costs 20 to 600 ms and under 1.5 GB, and anything that reads every column peaks at 5 to 15 GB. On the 78M x 43 tallyman entry, #1011 sorted a string column in 12 to 23 s and searched in 10 to 18 s, and its per-round peak RSS rose from 9.6 to 14.5 GB over three rounds.

Change

Adds docs/plans/ADR-004-scan-parquet-backend.md, status Proposed. It records nine decisions: a parquet /load on the polars backend opens a scan session; load costs a footer read; an unsorted window is a slice; a sorted window is a key phase over (row index, sort key) plus a gather of the window's rows; search separates the window from the count; order arrays and match lists are cached on disk by file identity; memory is a tested property; stats come from ADR-003 as a lazy select; and what a scan session does not offer. It has a measurement table, the alternatives rejected, and 8 open decisions, with a recommendation where the data supports one. Where a number was not measured the text says so.

No code changes.

Stack

Based on main. The ADR refers to ADR-001 (#1040), ADR-002 (#1043) and ADR-003 (#1044) by number; none of them has to merge first. The branch name has no slash for the Checks workflow's pull_request: branches: "*" trigger.

The probe script and its output are in the first comment.

🤖 Generated with Claude Code

Proposes a scan session for parquet /load on the polars backend: footer-only load, slice windows, a key phase plus gather for sorted windows, search with the window and count separated, on-disk order and match caches, a memory regression test, and stats through ADR-003. Includes measurements on a 12M x 43 file and the 78M #1011 run, and 8 open decisions.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
@chatgpt-codex-connector

Copy link
Copy Markdown

You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard.

@paddymul

paddymul commented Oct 6, 2026

Copy link
Copy Markdown
Collaborator Author

Probe scripts and raw output behind the "Measurements" table. Machine: Apple M4 Pro, load average 4.6, polars 1.35.2, one fresh process per case, one run each. The file is eval_993a_12m.parquet, the first 12M rows of parking_2017 (351 MB, 46 row groups, 43 columns, 29 of them String). Run as uv run python scan_probe.py CASE PATH from a checkout of main.

scan_probe.py

"""Window-cost probe for a scan backend. One case per process: uv run python scan_probe.py CASE PATH.

Prints one RESULT line: case, seconds, peak RSS (GB, from ru_maxrss), rows returned.
"""
import resource
import sys
import time

import polars as pl

case, path = sys.argv[1], sys.argv[2]
lf = pl.scan_parquet(path)
schema = lf.collect_schema()
strs = [c for c, d in schema.items() if d == pl.String]
SORT = "plate_id"  # String column
NUM = "summons_number"  # Int64 column
needle = "NY"


def mask():
    return pl.any_horizontal(*[pl.col(c).str.contains(needle, literal=True) for c in strs])


def two_pass(frame, col, desc, start, n, engine):
    keyed = frame.with_row_index("i")
    order = keyed.select("i", col).sort(col, descending=desc).slice(start, n).select("i").collect(engine=engine)
    rows = keyed.filter(pl.col("i").is_in(order["i"].implode())).collect(engine=engine)
    return order.join(rows, on="i", how="left", maintain_order="left")


t0 = time.perf_counter()
if case == "window_head":
    out = lf.slice(0, 100).collect(engine="streaming")
elif case == "window_deep":
    out = lf.slice(6_000_000, 100).collect(engine="streaming")
elif case == "count":
    out = lf.select(pl.len()).collect()
elif case == "sort_str_default":
    out = lf.sort(SORT).slice(100, 100).collect()
elif case == "sort_str_streaming":
    out = lf.sort(SORT).slice(100, 100).collect(engine="streaming")
elif case == "sort_str_two_pass":
    out = two_pass(lf, SORT, False, 100, 100, "streaming")
elif case == "sort_num_two_pass":
    out = two_pass(lf, NUM, False, 100, 100, "streaming")
elif case == "topk_str":
    out = lf.top_k(200, by=SORT, reverse=True).collect(engine="streaming")
elif case == "topk_num":
    out = lf.top_k(200, by=NUM, reverse=True).collect(engine="streaming")
elif case == "sort_idx_only_str":
    out = lf.with_row_index("i").select("i", SORT).sort(SORT).select("i").collect(engine="streaming")
elif case == "search_count":
    out = lf.filter(mask()).select(pl.len()).collect(engine="streaming")
elif case == "search_window_only":
    out = lf.filter(mask()).slice(0, 100).collect(engine="streaming")
elif case == "search_window_default":
    out = lf.filter(mask()).slice(0, 100).collect()
elif case == "search_rare_window_only":
    out = lf.filter(pl.col(SORT).str.contains("ZZZ9", literal=True)).slice(0, 100).collect(engine="streaming")
elif case == "topk_key_only_str":
    out = lf.with_row_index("i").select("i", SORT).top_k(200, by=SORT, reverse=True).collect(engine="streaming")
elif case == "topk_key_only_num":
    out = lf.with_row_index("i").select("i", NUM).top_k(200, by=NUM, reverse=True).collect(engine="streaming")
elif case == "gather_100_concat":
    import random
    random.seed(1)
    idx = sorted(random.sample(range(12_000_000), 100))
    out = pl.concat([lf.slice(i, 1) for i in idx]).collect(engine="streaming")
elif case == "gather_100_loop":
    import random
    random.seed(1)
    idx = sorted(random.sample(range(12_000_000), 100))
    out = pl.concat([lf.slice(i, 1).collect() for i in idx])
elif case == "search_chunked_1m":
    total = 0
    first = None
    for o in range(0, 12_000_000, 1_000_000):
        part = lf.slice(o, 1_000_000).filter(mask()).collect()
        total += part.height
        if first is None and part.height:
            first = part.head(100)
    out = first
    print("matches", total)
elif case == "select_str_cols_only":
    out = lf.select(strs).select(pl.len()).collect(engine="streaming")
else:
    raise SystemExit(f"unknown case {case}")
dt = time.perf_counter() - t0
rss = resource.getrusage(resource.RUSAGE_SELF).ru_maxrss / 1e9
n = out.height
print(f"RESULT {case:26s} {dt:7.2f} s  peak_rss {rss:5.2f} GB  rows {n}")

scan_probe_chunks.py

"""Per-chunk current RSS during a chunked search. uv run python scan_probe_chunks.py PATH CHUNK_ROWS"""
import os
import subprocess
import sys
import time

import polars as pl

path, chunk = sys.argv[1], int(sys.argv[2])
lf = pl.scan_parquet(path)
strs = [c for c, d in lf.collect_schema().items() if d == pl.String]
mask = pl.any_horizontal(*[pl.col(c).str.contains("NY", literal=True) for c in strs])


def rss_gb():
    return int(subprocess.check_output(["ps", "-o", "rss=", "-p", str(os.getpid())]).split()[0]) / 1e6


t0 = time.perf_counter()
total = 0
for k, o in enumerate(range(0, 12_000_000, chunk)):
    total += lf.slice(o, chunk).filter(mask).select(pl.len()).collect().item()
    if k % max(1, (12_000_000 // chunk) // 6) == 0:
        print(f"chunk {k} offset {o} rss_now {rss_gb():.2f} GB")
print(f"RESULT chunk_rows {chunk} matches {total} seconds {time.perf_counter() - t0:.2f} rss_end {rss_gb():.2f} GB")

Output of scan_probe.py, one process per case:

RESULT count                         0.01 s  peak_rss  0.07 GB  rows 1
RESULT window_head                   0.02 s  peak_rss  0.08 GB  rows 100
RESULT window_deep                   0.02 s  peak_rss  0.08 GB  rows 100
RESULT select_str_cols_only          0.03 s  peak_rss  0.34 GB  rows 1
RESULT topk_num                      1.13 s  peak_rss  6.08 GB  rows 200
RESULT topk_str                      0.41 s  peak_rss  7.45 GB  rows 200
RESULT sort_idx_only_str             0.60 s  peak_rss  1.22 GB  rows 12000000
RESULT sort_num_two_pass             0.53 s  peak_rss  7.52 GB  rows 100
RESULT sort_str_two_pass             1.00 s  peak_rss  7.52 GB  rows 100
RESULT search_rare_window_only       0.06 s  peak_rss  0.58 GB  rows 13
RESULT search_window_only            0.29 s  peak_rss  5.14 GB  rows 100
RESULT search_count                  0.41 s  peak_rss  5.42 GB  rows 1
RESULT topk_key_only_str             0.03 s  peak_rss  0.46 GB  rows 200
RESULT topk_key_only_num             0.02 s  peak_rss  0.26 GB  rows 200
RESULT gather_100_concat             0.47 s  peak_rss  7.50 GB  rows 100
RESULT gather_100_loop               1.24 s  peak_rss  0.34 GB  rows 100
RESULT search_chunked_1m             0.78 s  peak_rss  4.35 GB  rows 100
RESULT sort_str_streaming            3.03 s  peak_rss 10.51 GB  rows 100
RESULT sort_str_default              1.64 s  peak_rss 14.56 GB  rows 100
RESULT search_window_default         0.48 s  peak_rss  7.98 GB  rows 100

Output of scan_probe_chunks.py with 262,144 and then 1,000,000 rows per chunk:

chunk 0 offset 0 rss_now 0.31 GB
chunk 7 offset 1835008 rss_now 1.07 GB
chunk 14 offset 3670016 rss_now 1.73 GB
chunk 21 offset 5505024 rss_now 1.90 GB
chunk 28 offset 7340032 rss_now 1.92 GB
chunk 35 offset 9175040 rss_now 2.19 GB
chunk 42 offset 11010048 rss_now 2.21 GB
RESULT chunk_rows 262144 matches 10595588 seconds 1.59 rss_end 2.21 GB
chunk 0 offset 0 rss_now 0.98 GB
chunk 2 offset 2000000 rss_now 2.14 GB
chunk 4 offset 4000000 rss_now 2.75 GB
chunk 6 offset 6000000 rss_now 2.80 GB
chunk 8 offset 8000000 rss_now 2.84 GB
chunk 10 offset 10000000 rss_now 2.94 GB
RESULT chunk_rows 1000000 matches 10595588 seconds 0.96 rss_end 2.95 GB

search_chunked_1m in the first block is an earlier version of the same loop that kept the first 100 matches and reported peak RSS (4.35 GB); the second script reports current RSS at the end. Both use 1,000,000-row chunks. The "Measurements" table gives the second script's number for the chunked rows and says so.

The 78M-row figures come from the #1011 prototype (240a4b77) run against the tallyman entry, three rounds of six requests in one server process. The harness is not in the repo. Its log:

lazy_big_run1 rss before 0.163
load {'seconds': 2.429, 'resp': {'rows': 78024179}} rss 1.964
0 first_window_0_100 {'seconds': 0.005, 'length': 78024179, 'err': '', 'rows': 100, 'rss_after': 1.966, 'verify': 'ok'}
0 sort_str_100_200 {'seconds': 18.177, 'length': 78024179, 'err': '', 'rows': 100, 'rss_after': 5.424}
0 search_NY_0_100 {'seconds': 16.717, 'length': 68359359, 'err': '', 'rows': 100, 'rss_after': 8.651}
0 search_NY_5000_5100 {'seconds': 12.221, 'length': 68359359, 'err': '', 'rows': 100, 'rss_after': 9.562}
0 search_zero_match_0_100 {'seconds': 3.359, 'length': 0, 'err': '', 'rows': None, 'rss_after': 9.545}
0 sort_str_0_100 {'seconds': 14.968, 'length': 78024179, 'err': '', 'rows': 100, 'rss_after': 5.808}
1 first_window_0_100 {'seconds': 0.007, 'length': 78024179, 'err': '', 'rows': 100, 'rss_after': 5.812}
1 sort_str_100_200 {'seconds': 21.25, 'length': 78024179, 'err': '', 'rows': 100, 'rss_after': 8.788}
1 search_NY_0_100 {'seconds': 16.693, 'length': 68359359, 'err': '', 'rows': 100, 'rss_after': 8.406}
1 search_NY_5000_5100 {'seconds': 11.195, 'length': 68359359, 'err': '', 'rows': 100, 'rss_after': 11.14}
1 search_zero_match_0_100 {'seconds': 3.23, 'length': 0, 'err': '', 'rows': None, 'rss_after': 11.091}
1 sort_str_0_100 {'seconds': 12.328, 'length': 78024179, 'err': '', 'rows': 100, 'rss_after': 6.58}
2 first_window_0_100 {'seconds': 0.016, 'length': 78024179, 'err': '', 'rows': 100, 'rss_after': 6.584}
2 sort_str_100_200 {'seconds': 23.21, 'length': 78024179, 'err': '', 'rows': 100, 'rss_after': 9.248}
2 search_NY_0_100 {'seconds': 18.272, 'length': 68359359, 'err': '', 'rows': 100, 'rss_after': 9.639}
2 search_NY_5000_5100 {'seconds': 10.04, 'length': 68359359, 'err': '', 'rows': 100, 'rss_after': 14.454}
2 search_zero_match_0_100 {'seconds': 2.768, 'length': 0, 'err': '', 'rows': None, 'rss_after': 14.46}
2 sort_str_0_100 {'seconds': 12.673, 'length': 78024179, 'err': '', 'rows': 100, 'rss_after': 7.954}

@github-actions

github-actions Bot commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

📦 TestPyPI package published

pip install --index-strategy unsafe-best-match --index-url https://test.pypi.org/simple/ --extra-index-url https://pypi.org/simple/ buckaroo==0.15.9.dev37495732451

or with uv:

uv pip install --index-strategy unsafe-best-match --index-url https://test.pypi.org/simple/ --extra-index-url https://pypi.org/simple/ buckaroo==0.15.9.dev37495732451

MCP server for Claude Code

claude mcp add buckaroo-table -- uvx --from "buckaroo[mcp]==0.15.9.dev37495732451" --index-strategy unsafe-best-match --index-url https://test.pypi.org/simple/ --extra-index-url https://pypi.org/simple/ buckaroo-table

📖 Docs preview

🎨 Storybook preview

This branch was successfully deployed

1 active deployment
testpypi — 1ef0b254 Deployed Oct 6, 2026 by paddymul via Publish to TestPyPI #1756
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant