Repository navigation
Conversation
- 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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Closes #13.
What changed
results.parquet/part-<worker>.parquet, andmerge_parquet_filesis 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 existingresults.parquet(whether it's a file or a directory), so stale parts don't end up in the dataset.recv()and flushes only when a batch of 500 is full, or at shutdown.build_batchemits 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.WorkerConfig.output_vars, so every part has the same columns in the same order. This mattered foroutput_variables: ["all"]:runner.outputsis aHashMap, 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 foridare now set explicitly. Rows for each location are already contiguous, so row-group stats onidlet readers skip data when filtering.create_output_storenow takes&WorkerConfiginstead of separateglobal_start/first_locationarguments.Testing
cargo test: all tests pass, includingtest_all_formats_equivalent. The test reader now accepts either one file or a dataset directory, and checks that every part has the same schema.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.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.pd.read_parquet("results.parquet")on a 2-part directory gives the combined table, with zstd, dictionary-encodedid, and min/max stats.--no-default-features, and with onlyparquet. No new clippy warnings.Note
This changes the output location from a single
results.parquetfile to aresults.parquet/directory. It also edits the same region ofrun_parentas the NetCDF rework on theclaude/eager-goodall-eudoaobranch (which also makesdiscovered_varsunconditional) and as the Zarr PR #14 (which changeschunk_size). Whichever merges later may have small conflicts there.🤖 Generated with Claude Code
https://claude.ai/code/session_014HpTYAu1UyHLpFHz3ZVHcd