Repository navigation
Conversation
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>
|
You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard. |
|
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
"""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}")
"""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 Output of
The 78M-row figures come from the #1011 prototype ( |
📦 TestPyPI package publishedpip install --index-strategy unsafe-best-match --index-url https://test.pypi.org/simple/ --extra-index-url https://pypi.org/simple/ buckaroo==0.15.9.dev37495732451or 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.dev37495732451MCP server for Claude Codeclaude 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 |
Problem
The destination for polars
/loadis an out-of-corescan_parquetbackend (#1000 and #1005, the eager stopgaps, are closed). No design for it is in the repo: ADR-003 lists "thescan_parquetbackend (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/loadon 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'spull_request: branches: "*"trigger.The probe script and its output are in the first comment.
🤖 Generated with Claude Code