Sitelet https://github.com/buckaroo-data/buckaroo/pull/1012
Skip to content

perf(stats): batch polars value_counts across columns with collect_all (#997) - #1012

Closed
paddymul wants to merge 2 commits into
mainfrom
perf/997-batch-value-counts
Closed

paddymul wants to merge 2 commits into
mainfrom
perf/997-batch-value-counts

Conversation

@paddymul

@paddymul paddymul commented Oct 3, 2026 •

Copy link
Copy Markdown
Collaborator

Part of #997. One of several candidate approaches; the alternatives are taking mode from the first row of value_counts only (perf/997-mode-from-value-counts) and the batch pipeline / stats pre-pass work under #999.

Problem

pl_base_summary_stats runs value_counts and mode on one pl.Series at a time, so polars never parallelises them across columns. On the 10.8M-row, 43-column parquet from #997 that is 6.8 s of value_counts plus 2.5 s of mode, against 2.8 s when the same 43 value_counts are collected together with pl.collect_all.

Approach

PlDfStatsV2 computes value_counts for every column in one pl.collect_all before process_df runs, converts each result to the pd.Series shape the rest of the DAG expects, and hands them to StatPipeline.process_df as per-column initial stats. pl_base_summary_stats takes value_counts from the accumulator and derives mode from its first row instead of running a second group-by. A per-series provider, pl_value_counts, stays in PL_ANALYSIS_V2 for callers that use StatPipeline directly; when a column's value_counts is supplied up front the pipeline skips that stat, so the key is only ever provided once.

What changes

  • StatPipeline.process_df accepts column_initial_stats (orig column name to stats dict), merged into each column's initial_stats.
  • StatPipeline.process_column skips a stat whose every provided key was supplied through initial_stats.
  • polars_utils.batch_value_counts(df, columns) runs the one-pass collect and returns {column: pd.Series}; vc_frame_to_pd is the shared conversion used by both paths. Object columns are left out (polars panics on them inside collect_all).
  • pl_stats_v2: pl_value_counts provides value_counts; pl_base_summary_stats requires it and derives mode from it.
  • PlDfStatsV2 runs the batch when the pipeline provides value_counts, falls back to the per-series path if the collect raises, and reuses the batch in add_analysis.

The 50k-row sample in PlDfStatsV2.get_operating_df is unchanged; #992 covers it.

Tests

tests/unit/test_pl_stats_v2.py::TestPlBatchValueCounts: pl.Series.value_counts is monkeypatched with a counter and PlDfStatsV2 on a 5-column frame makes zero per-Series calls; the batched pipeline output matches the per-series output key by key on a frame with no count ties; a supplied value_counts pre-empts the per-series stat (mode and distinct_count follow it); process_column with a raw series and no initial value still produces value_counts and mode; Object columns are left out of the batch and still get stats.

Compatibility

pl_base_summary_stats now requires value_counts, and the provider of that key moved to the new pl_value_counts stat. A pipeline built by hand from the old public name, for example a list copied from an older PL_ANALYSIS_V2 such as [pl_typing_stats, _type, pl_base_summary_stats, ...], raises DAGConfigError: No function provides 'value_counts' at construction. Adding pl_value_counts to the list fixes it. In-repo callers use PL_ANALYSIS_V2 and are unaffected, and a custom stat that provides both value_counts and mode still works. The existing tests that built StatPipeline([pl_base_summary_stats]) directly were updated to include pl_value_counts.

Trade-offs

  • mode among tied values was unspecified before and still is; it now follows value_counts order, so mode and most_freq always agree.
  • mode for datetime and duration columns comes back as pd.Timestamp / pd.Timedelta (subclasses of the stdlib types, and what most_freq already returned) rather than datetime / timedelta.
  • If collect_all raises for any column the whole batch is dropped and every column goes through the per-series path, so a frame with one unsupported column pays the old cost plus a failed attempt.
  • A custom provider of value_counts passed to PlDfStatsV2 is pre-empted by the batch; StatPipeline used directly is unaffected.

🤖 Generated with Claude Code

paddymul and others added 2 commits October 2, 2026 20:07
PlDfStatsV2 should compute value_counts for every column in one
pl.collect_all instead of one Series at a time, hand each column its
result through StatPipeline.process_df, and derive mode from it.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
#997)

PlDfStatsV2 computes value_counts for every column in one pl.collect_all
before the per-column pass and seeds each column's accumulator with it
through the new StatPipeline.process_df(column_initial_stats=...).
process_column skips a stat whose every provided key was seeded, so the
batch is the only provider of value_counts on that path.

pl_base_summary_stats is split: pl_value_counts provides value_counts
per Series (for StatPipeline callers that pass a raw series), and
pl_base_summary_stats takes value_counts from the accumulator and
derives mode from its first row instead of a second group-by.

Object columns are left out of the batch because collect_all panics on
them; a batch that raises falls back to the per-Series path.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
@github-actions

github-actions Bot commented Oct 3, 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.dev37085735205

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.dev37085735205

MCP server for Claude Code

claude mcp add buckaroo-table -- uvx --from "buckaroo[mcp]==0.15.9.dev37085735205" --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

@paddymul

paddymul commented Oct 3, 2026

Copy link
Copy Markdown
Collaborator Author

Review of this PR found two issues, both addressed.

  1. TDD sequence incomplete. The tests-only commit d146997 had a failing Checks run, but the fix commit 948902d had not been pushed, so the PR had no passing run. I pushed 948902d; Checks run 37085735205 passes on every job (Python 3.11-3.14, Windows, max-versions, lint, JS, wheel, Jupyter/Marimo/Server Playwright, MCP). The deploy job is skipped. The reviewer also noted that the test commit fails at collection for the whole of test_pl_stats_v2.py because it imports the new API at module level, not just the new tests. I left the history as it is: the failure was seen on CI, and rewriting pushed commits to split the imports would discard that run.

  2. Public API break. pl_base_summary_stats now requires value_counts, provided by the new pl_value_counts stat, so a hand-built pipeline using the old name without pl_value_counts raises DAGConfigError at construction. I did not add a back-compat path: keeping pl_base_summary_stats self-sufficient would put a second provider of value_counts and mode back into the same stat, which is what the batching removes, and in-repo callers all go through PL_ANALYSIS_V2. Instead I added a Compatibility section to the PR body that names the break, shows the error, and gives the fix (add pl_value_counts to the list; I checked that [pl_typing_stats, pl_value_counts, pl_base_summary_stats] builds). No code change, so no new commit for this item.

Local run before pushing: ruff clean, unit suite 1174 passed. Fifteen lazy-widget tests failed locally with sqlite3.OperationalError: database is locked on the shared ~/.buckaroo/executor_log.sqlite, because other worktrees were running tests at the same time. Re-running those three files with a private HOME gave 25 passed, and CI is green.

@paddymul

paddymul commented Oct 3, 2026

Copy link
Copy Markdown
Collaborator Author

Measurements

Timed create_polars_dataflow on the branch head (948902d) against main (992fdb3), each in a fresh uv run python process, 3 processes per cell (first call in the process; two further in-process calls gave the same numbers within 0.02 s on the small file and 0.15 s on the large one). "Full" means the 50k sample is off (pre_limit = False and get_operating_df returning the frame unchanged), as in bench_fullstats.py. "Default" leaves the sampling as shipped. Machine load average was about 4 throughout. The 78M-row file was not used.

Case main (median) branch (median) Change
parking_2017.parquet, 10.8M x 43, full 8.45 s (8.31-8.50) 4.14 s (4.06-4.15) -51%
8e98b80a36ef.parquet, 52,814 x 27, full 0.331 s 0.257 s -22%
8e98b80a36ef.parquet, 52,814 x 27, default (sample on) 0.410 s 0.343 s -16%

Peak RSS (resource.getrusage, one create_polars_dataflow call on the 10.8M-row file, full): main 7.7 GB, branch 10.6-11.1 GB. With three calls in one process it was 8.2 GB vs 13.2 GB. Both are under the 16 GB limit, but the branch holds about 3 GB more at peak on that file, which is the value_counts frames of all 43 columns being materialised by one collect_all rather than one column at a time. On the small file RSS was within 1.0-1.9 GB for both and varied from process to process by more than the difference between branches.

The search-rerun case on the small file measured 0.41 s on main here, not the ~588 ms expected; I did not find what accounts for the difference.

Correctness

Compared summary_sd on 8e98b80a36ef.parquet (27 columns, 1,122 column/key pairs) between main and the branch, with full stats. Two runs of main against each other already differ in the mode of __row_order and player_page: both columns have ties at the maximum count (52,814 values at count 1, and 2 values at count 37), and Series.mode() does not order ties. Against main, the branch differs in five mode entries and nothing else:

  • __row_order, player_page, otc_id: same tie situation. The branch now returns the first row of value_counts, which equals most_freq in all three cases. Main's mode disagreed with its own most_freq here.
  • season_history, contract_history (list-of-struct columns): mode is a numpy array instead of a pl.Series, equal in content to most_freq.

Every other key, including most_freq, 2nd_freq to 5th_freq, histogram_args, distinct_count and value_counts, has an identical repr.

Summary: the large-file stats phase is about twice as fast, at the cost of about 3 GB more peak memory on a 43-column, 10.8M-row frame. The only output differences are mode on tied columns and on list columns.

@paddymul

paddymul commented Oct 5, 2026

Copy link
Copy Markdown
Collaborator Author

Closing. Batching value_counts across columns with collect_all tunes the in-memory path and costs about 3 GB more peak memory at 10.8M rows; stats for the out-of-core backend come from a single lazy select. The branch stays.

@paddymul paddymul closed this Oct 5, 2026
@paddymul paddymul added the buckaroo-server-polars Polars /load work on the buckaroo server (#992-#999), superseded by a scan_parquet backend label Oct 5, 2026

This branch was successfully deployed

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

Labels

buckaroo-server-polars Polars /load work on the buckaroo server (#992-#999), superseded by a scan_parquet backend

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant