Skip to content

Zarr output: multi-location chunks, buffered writes, zstd - #14

Open
JoshCu wants to merge 1 commit into
mainfrom
claude/relaxed-goldberg-qmrkqi
Open

JoshCu wants to merge 1 commit into
mainfrom
claude/relaxed-goldberg-qmrkqi

Conversation

@JoshCu

@JoshCu JoshCu commented Sep 26, 2026

Copy link
Copy Markdown
Owner

Closes #12.

What changed

  • Chunks hold many locations. Output variables are now chunked [B, n_times] instead of [1, n_times]. zarr::chunk_locations(total_steps, per_worker) sizes B for about 8 MB of f32 per chunk, capped at 1024 and at one worker's share of the locations. run_parent rounds each worker's chunk_size up to a multiple of B, so every chunk belongs to exactly one worker and concurrent writes stay lock-free.
  • Buffered writes on a background thread. ZarrStore opens the store and every array once. It then buffers locations on a writer thread, like the other writers, and writes one store_chunk per variable per chunk. Short series and missing variables are padded with NaN, and so are the rows past the end of the last partial chunk. If a block doesn't line up with a chunk (a caller-chosen start_idx), it falls back to store_array_subset.
  • IDs: all IDs are written in one call, into a single chunk.
  • zstd (level 3) on the variable and id arrays.
  • The root group gets a variables attribute that lists the output arrays, so workers can open them all up front.
  • create_zarr_store takes a new chunk_locations argument, and ZarrStore::new now returns BmiResult.

Files

For 7 locations with B = 3, each variable now has 3 chunk files (c/0/0, c/1/0, c/2/0) instead of 7. In general it's ceil(n_locations / B) per variable.

Testing

  • cargo test: all tests pass, including test_all_formats_equivalent.
  • New test test_zarr_multi_worker_partial_chunks: 23 locations, 4 workers and B = 5, so 3 aligned workers write concurrently from threads, and the last chunk is partial. One location has a short series and another is missing a variable. The test checks the values, the NaN fill, and that each variable has ceil(23/5) = 5 chunk files.
  • Checked by hand that a store written this way opens in xarray (xr.open_zarr, zarr-python 3.1) with correct values and ZstdCodec(level=3).
  • Builds with default features, with --no-default-features, and with only zarr. No new clippy warnings.
  • Not run: a full end-to-end model run, since that needs compiled model libraries.

Note

This touches the Zarr block of run_parent, which the NetCDF rework (the claude/eager-goodall-eudoao branch) also edits, so whichever merges second may have a small conflict there.

🤖 Generated with Claude Code

https://claude.ai/code/session_014HpTYAu1UyHLpFHz3ZVHcd


Generated by Claude Code

- Chunk output variables as [B, n_times] instead of [1, n_times], with B
  sized for ~8 MB per chunk (max 1024, and at most one worker's share).
  run_parent rounds each worker's block up to a multiple of B, so every
  chunk belongs to one worker and concurrent writes stay lock-free.
- ZarrStore opens the store and arrays once, and buffers locations on a
  background thread, writing one store_chunk per variable per chunk.
  Short series and missing variables are padded with NaN.
- Write all ids in one call into a single chunk.
- zstd-compress the variable and id arrays.
- Add a test with several concurrent workers, a location count that is
  a multiple of neither B nor the worker count, a short series and a
  missing variable, and check the chunk file count.

Closes #12

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.

Zarr output: one chunk per catchment, no batching, store reopened on every write

2 participants