Skip to content

Parquet output: write a dataset directory, drop the re-encoding merge - #15

Open
JoshCu wants to merge 1 commit into
mainfrom
claude/parquet-dataset-13
Open

JoshCu wants to merge 1 commit into
mainfrom
claude/parquet-dataset-13

Conversation

@JoshCu

@JoshCu JoshCu commented Sep 26, 2026 •

Copy link
Copy Markdown
Owner

Closes #13.

What changed

  • No merge (option A from the issue). Each worker writes results.parquet/part-<worker>.parquet, and merge_parquet_files is gone. Nothing rewrites the data single-threaded after the workers finish. pandas, polars, pyarrow, DuckDB and Spark all read the directory as one table. Before spawning workers, the parent removes any existing results.parquet (whether it's a file or a directory), so stale parts don't end up in the dataset.
  • No timer flush. The writer thread blocks on recv() and flushes only when a batch of 500 is full, or at shutdown.
  • No panic on short columns. build_batch emits rows up to the longest series and fills short or missing ones with NaN (the same fill value as Zarr and NetCDF). It also looks each variable up once per location instead of once per row.
  • Fixed schema. The parent's discovered variable list is no longer zarr-only. It goes to workers in WorkerConfig.output_vars, so every part has the same columns in the same order. This mattered for output_variables: ["all"]: runner.outputs is a HashMap, so column order from the first result could differ between workers. If the list is empty, the writer falls back to the first result's columns.
  • max_row_group_row_count (1M rows) and dictionary encoding for id are now set explicitly. Rows for each location are already contiguous, so row-group stats on id let readers skip data when filtering.
  • To stay under clippy's argument limit, create_output_store now takes &WorkerConfig instead of separate global_start / first_location arguments.

Testing

  • cargo test: all tests pass, including test_all_formats_equivalent. The test reader now accepts either one file or a dataset directory, and checks that every part has the same schema.
  • New test_parquet_short_and_missing_columns: one variable is shorter than the first column, another is missing, and a second location has its columns in a different order and a longer last column. The test checks there's no panic and that the NaN fill is correct.
  • New test_parquet_dataset_parts: two writers on separate threads write parts of one directory, each with its columns in reverse order, and the directory is read back as one table.
  • Checked by hand that pd.read_parquet("results.parquet") on a 2-part directory gives the combined table, with zstd, dictionary-encoded id, and min/max stats.
  • Builds with default features, with --no-default-features, and with only parquet. No new clippy warnings.
  • Not run: a full end-to-end model run, since that needs compiled model libraries.

Note

This changes the output location from a single results.parquet file to a results.parquet/ directory. It also edits the same region of run_parent as the NetCDF rework on the claude/eager-goodall-eudoao branch (which also makes discovered_vars unconditional) and as the Zarr PR #14 (which changes chunk_size). Whichever merges later may have small conflicts there.

🤖 Generated with Claude Code

https://claude.ai/code/session_014HpTYAu1UyHLpFHz3ZVHcd

- Workers write results.parquet/part-<worker>.parquet; the parent no
  longer decodes and re-encodes every worker file after the run.
- The writer thread blocks on recv() and flushes only when a batch is
  full or at shutdown, instead of every 100 ms.
- The schema comes from the variable list the parent discovers, passed
  to workers in WorkerConfig, so every part has the same columns in the
  same order (runner.outputs is a HashMap, so first-result order varied).
- build_batch emits rows up to the longest series and NaN-fills short
  or missing ones instead of indexing past the end.
- Set max row group size and dictionary encoding for id explicitly.
- Tests for a short / missing variable and for multi-part datasets.

Closes #13

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014HpTYAu1UyHLpFHz3ZVHcd
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.

Parquet output: merge decodes and re-encodes everything; timer flush; panic on short columns

2 participants