Skip to content

DDStore 3.0: GPUDirect RDMA, batched reads (get_batch), thread-safe get() - #6

Merged
jychoi-hpc merged 56 commits into
mainfrom
check-thread
Oct 5, 2026
Merged

jychoi-hpc merged 56 commits into
mainfrom
check-thread

Conversation

@jychoi-hpc

Copy link
Copy Markdown
Member

DDStore 3.0: GPUDirect RDMA, batched reads, and thread-safe get()

This brings the GPUDirect work and the follow-up threading/performance work from check-thread into main, and sets the package version to 3.0 (v2.0 is the previous release tag). Suggested merge: squash (the branch history includes early wip/add commits); tag v3.0 on the merge commit.

What's new

GPU-resident buffers (GPUDirect RDMA, method=1/2, DDSTORE_FABRIC=cxi)

  • add() and get() accept a CUDA/HIP torch.Tensor: RDMA reads from / writes into GPU memory with no host copy. hsn and method=0 reject GPU buffers with a clear error.
  • get() syncs the device before each GPU-destination transfer; removing it was shown to cause GPU memory faults.

Batched reads (default)

  • PyDDStore.get_batch(name, arr, indices): one call per training batch.
    • method=1/2: one lock, one memory registration, one GPU sync, and all of the batch's fi_reads in flight together.
    • method=0: collective, after MDLoader (IPDPSW 2024): MPI_Allgatherv of indices + MPI_Alltoallv of rows, in bounded rounds (DDSTORE_ALLTOALL_MAX_BYTES, default 2 MiB).
  • DistDataset/DistDatasetReader.__getitems__ use it, so PyTorch's DataLoader and ThreadDataLoader read whole batches automatically; DDSTORE_BATCH_GET=0 switches back to per-sample reads.

Threading

  • get() is thread-safe: a per-variable lock inside DDStore::get() (needed with 2+ threads; uncontended otherwise). get() releases the GIL on both host and GPU paths.
  • ThreadDataLoader (thread-based workers, safe with GPU buffers) collates in the worker thread.
  • mpi4py.rc settings removed: only the main thread calls MPI.

Fixes

  • free() releases the host buffers DDStore allocated (was leaking MPI_Alloc_mem); join() initializes its mutex.
  • job-vae-core-extra.sh (method 2 core/extra as separate srun steps) now works on Slingshot: --network=single_node_vni,job_vni plus keeping only the job VNI in SLINGSHOT_VNIS. Defaults to --layout=split-node; colocate works on Perlmutter only (with --overlap).
  • Stale comments/README sections corrected; flaky test teardown fixed.

Tools

  • DDSTORE_PROFILE=1 + get_profile(): where get()/get_batch() time goes.
  • examples/scripts/bench_get.py: per-row latency/throughput vs row size, destination, batch size, threads.
  • VAE example: --replicate R, --image-scale S; test/test_get_batch.py; examples/scripts/perlmutter-check.sh.
  • README rewritten around usage (environment-variable reference, GPUDirect, concurrency, Slingshot multi-step notes); all measurements in docs/results.md; MDLoader added to the citations.

Results

Full tables and the experiments behind them: docs/results.md.

VAE (vae-ddp.py), seconds per epoch, per-sample → batched, identical losses in every pair:

Frontier, 4 nodes × 8 Perlmutter, 2 nodes × 4
method 1, host, 0 workers 0.128 → 0.088 0.398 → 0.208
method 1, GPU, 2 workers (2 nodes) 0.403 → 0.134 0.678 → 0.183
method 0, host 0.241 → 0.088 0.849 → 0.219

Per-row single-get cost was dominated by the GPU sync waiting on training kernels (130–430 µs/sample with workers); batching brings it to ~0.1 µs/row. GPUDirect vs host destination is platform-dependent: GPU wins from ~12.5 KB rows on Frontier, host wins at all sizes on Perlmutter.

Testing

  • Frontier (MI250X, ROCm 7.2, cxi): test_single, test_multirank, test_gpu_rdma 14/14, test_get_batch 20/20 on 4 ranks and 19/19 on 16; VAE losses unchanged across all variants; core/extra split-node (host and --gpudirect).
  • Frontier, DDSTORE_FABRIC=hsn: test_get_batch 19 passed / 1 skipped (cxi-only GPU case), test_gpu_rdma 14/14, VAE method 1, GPU buffers refused with a clear error, core/extra split.
  • Perlmutter (A100, CUDA 13, cxi): all suites pass incl. CUDA GPUDirect; VAE identical losses; split-node and colocate (--overlap).

Known limitations

  • Colocated core/extra steps fail to launch on Frontier (Error configuring interconnect), with or without --overlap; use split-node.
  • Host-destination batches of 1 MB rows on Frontier are slower than single reads (not investigated).
  • Job scripts' #SBATCH lines target Frontier (-A FUS184, 8 ranks × 7 cores per node).

🤖 Generated with Claude Code

https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z

jychoi-hpc and others added 30 commits September 3, 2026 11:02
Direct Frontier experiments confirmed the threading.Lock in
DistDataset/DistDatasetReader.get() is necessary: disabling it crashed
with "double free or corruption" under concurrent thread access (the
libfabric CQ/MR-cache state in src/common.cxx has no locking of its own).

Found a separate, deeper bug while verifying: ThreadDataLoader's GPU
buffer pool hands out slots per-sample via lock-acquisition order, not
batch order, so bounding in-flight batch count (implemented here, with a
refill-ordering fix and pool_size sizing) reduces but cannot fully
eliminate silent data corruption with --num-workers > 1 together with
--gpu-dest/--gpu-source. vae-ddp.py/vae_extra_train.py now raise a clear
RuntimeError for that combination instead of running silently wrong.

Kept examples/vae/stress_threaded_loader.py (the diagnostic harness used
to find both issues) and documented the root cause and what a complete
fix needs in the README for a later revisit.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01443BqKw3NHnbGg8HSeoLXJ
…dings

- pyddstore.pyx: add() and get() normalize GPU-tensor / numpy input to
  (ptr, itemsize, iface) and make one C++ call dispatched on item size
  (DDStore::add/get<T> only use T via sizeof). Shared _check_dtype keeps
  the supported-dtype list. get() now releases the GIL on the GPU path too.
  The per-call torch.cuda.synchronize() in get() is kept: removing it made
  every --gpu-dest VAE run abort with HSA_STATUS_ERROR_EXCEPTION.
- Remove mpi4py.rc.thread_level/threads from examples, tests and README:
  only the main thread calls MPI; SINGLE/FUNNELED/MULTIPLE gave identical
  results and timing on Frontier.
- test_gpu_rdma.py: barrier before free() in the compute-kernel-read test
  (fast rank freed its endpoint while the peer was still reading ->
  PTLTE_NOT_FOUND, intermittent before this change too).
- README: no-build-isolation install, lock necessity, GIL behaviour,
  worker-count timings, MPI thread level, HIP hardware-queue limits.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
fetch() now runs collate_fn (then pin_memory) in the worker instead of
__next__ doing it on the training thread. Moves the per-batch stack off
the training loop and makes pin_memory=True effective (previously it
pinned per-sample tensors that collate then copied into an unpinned one).

Frontier, vae-ddp.py method=1 cxi, 16 ranks, 1 worker, host path:
fetch ~0.05 -> ~0.006 s/epoch, epoch 0.20 -> 0.15 s. Loss unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
… mutex

Docs (comments, docstrings, help text, README):
- cxi is the native Slingshot provider on Frontier and Perlmutter (required
  for GPUDirect), not "Perlmutter only"; drop old branch-name references.
- Thread-safety note: only hsn requests FI_THREAD_DOMAIN; the lock is
  needed for the shared per-variable recv fields.
- Remove descriptions of the removed GPU buffer pool (recv-MR cache).
- Fix update() signature, epoch no-op scope, demo/test.py backend, script
  --method/--num-workers help, get_local_rank docstring, stale test notes,
  missing script references; add post-collate-change timings.

Code:
- free(): track owns_base per variable and MPI_Free_mem the host buffer
  DDStore allocated in add()/init(), after closing the window/MRs that
  reference it; skip MPI calls after MPI_Finalize; destroy recv_lock.
- join(): pthread_mutex_init the fabric_state lock like every other
  allocation site (get() locks it on the extra side too).
- ThreadDataLoader: drop unused _dataset_fetcher.

Tests (Frontier, cxi): test_single (method 0/1) 14/14, test_multirank 5/5,
test_gpu_rdma 14/14.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
…kers

On Slingshot every srun step gets its own VNI, and endpoints only talk on
the same VNI. The core/extra split therefore needs:
- #SBATCH --network=single_node_vni,job_vni: job_vni adds a job-wide VNI
  to each step (SLINGSHOT_VNIS=<step VNI>,<job VNI>); single_node_vni gives
  single-node steps a CXI service at all (else fi_domain fails with -38).
- a per-task wrapper keeping only the job VNI (last entry), since the cxi
  provider uses only the first VNI listed (else VNI_NOT_FOUND).

Verified on Frontier, method=2, cxi, split-node: 1 core + 1 extra node and
2 core + 4 extra nodes both train 3 epochs and shut down cleanly; removing
either piece reproduces the failure. Colocate (both steps on the same
nodes) still fails to launch the second step with job_vni ("Error
configuring interconnect"); without job_vni the steps can't share a VNI.

Also add --num-workers (extra step) and document the fix in README.

Correction to 80c3c65's message: test_single/test_multirank always use
method=0 (DDSTORE_METHOD is ignored), so its "method 0/1" results were
method 0 only; methods 1/2 were covered by test_gpu_rdma and VAE runs.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
- DDSTORE_PROFILE=1: per-variable C++ counters in fabric_state (calls,
  lock wait, recv-MR check/registration + misses, fi_read post, CQ wait),
  updated under recv_lock; DDStore::profile() and PyDDStore.get_profile(),
  which adds Python-side whole-call and GPU sync time. Off by default.
- vae-ddp.py prints an all-rank get() breakdown at the end when enabled.
- examples/scripts/bench_get.py: single-row get() latency/throughput vs
  row size, host / fresh GPU / reused GPU destination, threads, with the
  same breakdown.
- VAE: --replicate R (ConcatDataset, longer epochs) and --image-scale S
  ((28*S)^2 images, hidden 400*S); core server and both job scripts pass
  them through, extra side follows the published row width. Defaults
  reproduce the original example (epoch-8 loss 8.8960 at R=1/S=1).
- README: new options, profiling, and Frontier measurements.

Findings (Frontier, cxi, 2x8 ranks): GPU-destination get() inside the VAE
is dominated by the per-call device sync once a worker thread runs
(130-430 us/sample vs ~4 us without workers); lock wait and MR
registration are negligible. GPUDirect beats host between 12.5 KB and
200 KB per row (1 MB: 134 vs 206-276 us).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
Step-by-step plan with the measured per-sample baseline (DDSTORE_PROFILE,
bench_get.py) it is meant to beat.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
One call per training batch instead of one get() per sample.

- C++: read_from_remote() split into ensure_recv_mr(), reap_one() and
  read_batch_from_remote(), which posts all of a batch's fi_read()s
  (reaping completions when the TX queue is full) before waiting, and
  always waits for every posted read, even after an error, so no stale
  completion is left for the next call. Single-row get() is a 1-row batch.
  DDStore::get_batch() validates every index first, then takes the
  per-variable lock once; method 0 loops the per-row MPI_Get path.
  DDSTORE_PROFILE also counts rows.
- Cython: PyDDStore.get_batch(name, arr, indices): one GPU sync per batch,
  GIL released, ValueError on rows/indices mismatch.
- DistDataset / DistDatasetReader.__getitems__ (used by DataLoader and
  ThreadDataLoader); DDSTORE_BATCH_GET=0 falls back to per-sample get().
- bench_get.py --batch N (per-row reporting); vae-ddp profile summary per
  call and per row.
- test/test_get_batch.py: method 0 and method 1 (cxi), GPU destination,
  concurrent threads, error recovery.

Frontier (cxi, 2 nodes x 8 ranks): new tests 20/20 on 4 ranks; existing
suites pass. VAE epoch time with batching, S=1: host w1 0.156 -> 0.126 s,
GPU w1 0.259 -> 0.127 s, GPU w2 0.403 -> 0.134 s (S=2: GPU w1 0.416 ->
0.264 s); losses unchanged (8.8960 / 30.0178). Per-row GPU sync drops from
129-383 us to ~0.1 us. bench: 3 KB rows 9.0 -> 0.68 us/row (host),
20.9 -> 0.78 (GPU); GPUDirect now beats host from 12.5 KB rows up.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
For method 0, DDStore::get_batch() now all-gathers every rank's batch
indices, packs the rows each rank owns per requester, and delivers them with
one MPI_Alltoallv (rows as a contiguous datatype), instead of a per-row
MPI_Win_lock/MPI_Get/unlock loop. Runs on a private duplicate of the
store's communicator (freed in free()); indices are validated on the
gathered list so every rank raises together. get_batch() is therefore
collective for method 0: every rank calls it the same number of times, in
the same order, from one thread at a time. __getitems__, loaders, methods
1/2 and single-row get() are unchanged.

Also: README environment-variable section and get_batch notes; bench_get.py
per-row normalization for method 0 (no C++ counters there).

Frontier (2 nodes): test_get_batch 20/20 on 4 ranks, 19/19 on 16; other
suites pass. VAE method 0, per-row -> collective: S=1 0.421 -> 0.142 s/epoch,
S=2 0.597 -> 0.315 (losses unchanged). bench (16 ranks, us/row): 3 KB
30 -> 2.95, 12.5 KB 35 -> 7.4; large rows lose at batch 128 (200 KB 272,
1 MB 1765 vs 166 / 723 per-row) - to be addressed by bounded rounds.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
…BYTES)

One exchange of a whole batch of large rows was slower than per-row reads
(e.g. 1 MB rows at batch 128: 1765 vs 739 us/row). Split it into rounds of
at most DDSTORE_ALLTOALL_MAX_BYTES received per rank (default 2 MiB); every
rank derives the same round count from the gathered request counts.

Frontier, 16 ranks, us/row at batch 128 (per-row get -> no cap -> 2 MiB):
200 KB 170 -> 272 -> 111, 1 MB 739 -> 1765 -> 685; 3 KB / 12.5 KB unchanged
(one round). 32 MiB reproduces the slowdown; 2 and 8 MiB are similar.
Tests pass on 4 ranks and on 16 ranks with a 50-byte cap (many rounds);
VAE losses unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
jychoi-hpc and others added 26 commits October 4, 2026 16:03
- test_get_batch.py: method-1 cases use DDSTORE_FABRIC when set (cxi
  otherwise); the GPU-destination case runs on cxi only. On Frontier with
  DDSTORE_FABRIC=hsn: 19 passed, 1 skipped (4 ranks); test_gpu_rdma 14/14;
  VAE method 1 over hsn, batched on/off x 0/2 workers, loss 8.8960; GPU
  destination refused with the hsn error; method-2 split over hsn trains.
- examples/scripts/perlmutter-check.sh + docs/perlmutter-checklist.md:
  what Frontier could not cover (CUDA GPUDirect, Perlmutter Slingshot VNIs,
  colocate), as one 2-node debug job writing pm-check-<jobid>/summary.txt.
  Syntax-checked only; not yet run on Perlmutter.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
Perlmutter check (job 59332134, 2 nodes x 4 A100): all tests pass incl.
CUDA GPUDirect (test_gpu_rdma 14/14, test_get_batch 20/20 on 8 ranks);
VAE losses identical for every variant, batched 1.9-3.9x faster per epoch;
core/extra split-node passes; VNI order <step>,<job> as on Frontier.

- Colocate works on Perlmutter only with job_vni + the VNI wrapper + srun
  --overlap; job-vae-core-extra.sh now runs colocate steps with --overlap.
  On Frontier the same setup still fails ("Error configuring interconnect",
  re-tested with --overlap): use split-node there.
- Perlmutter single-node steps work without single_node_vni (default CXI
  service); Frontier needs it.
- GPUDirect vs host crossover is platform-dependent: host destinations win
  at every row size on Perlmutter (1 MB: 75 vs 142 us/row at batch 128).
- perlmutter-check.sh: match rank 0 with or without srun -l padding (loss
  columns were empty with < 10 ranks); colocate case uses --overlap.
- README and checklist updated with the above.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
Colocate fails on Frontier (the script's #SBATCH defaults target it) and
works on Perlmutter only with --overlap; split-node works on both.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
…cite MDLoader

- README second half reorganized around usage: GPUDirect (rules, not the
  investigation), PyTorch integration (batched by default, recommended
  settings), job scripts, sub-communicators, a short Performance section,
  Known Limitations (multi-step Slingshot + colocate table, Concurrency,
  troubleshooting). Removes the stale "GPU-to-GPU performance" section.
- docs/results.md: all measurements and experiments (batched vs per-sample
  on Frontier 2/4 nodes and Perlmutter, DDSTORE_PROFILE breakdown,
  bench_get.py, method-0 collective and round sizes, worker threads, GPU
  sync / lock / pool / MPI thread level experiments, HIP streams).
- Citation: add MDLoader (SC24-W, doi 10.1109/SCW63240.2024.00145), which
  the method-0 batched path follows.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
v2.0 is already tagged (2025-09-18, PR #4); this release adds GPUDirect
RDMA, batched reads (get_batch) and thread-safe get(), and changes
behavior (collective method-0 get_batch, free() releases memory,
core/extra script default layout).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
Pick the rank layout from NERSC_HOST/LMOD_SYSTEM_NAME. On Perlmutter, DDP
ranks see all 4 GPUs of their node (--gpus-per-node=4) and pick
cuda:$SLURM_LOCALID: with --gpus-per-task=1, NCCL 2.29 (pytorch/2.13.0)
fails in DDP init with "Cuda failure 101 'invalid device ordinal'".
Frontier keeps 8 ranks/node with --gpus-per-task=1.

job-vae-single.sh: add --ranks-per-node and --cpus-per-task.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016wo6ZbZfw9fvsZooyVL3j2
…ze=None

- clean(): drain the future queue without blocking. The end marker is
  queued only once the sampler is exhausted, so after an epoch that
  stopped early the old iter(fs.get, None) blocked forever.
- __iter__(): draw a base seed from the generator every epoch, as torch's
  DataLoader iterator does, so later random draws match DataLoader.
- fetch(): with batch_size=None (no auto-collation) pass the sampler's
  index to dataset[index] as is, as torch's map-style fetcher does.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016wo6ZbZfw9fvsZooyVL3j2
Match 2d5017e: with --gpus-per-task=1, NCCL 2.29 fails in DDP init on
Perlmutter ("invalid device ordinal"). VAE runs and the extra (training)
step of the core/extra checks now use --gpus-per-node=4; pytest and
bench_get.py steps keep --gpus-per-task=1 (no NCCL).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
…ts.md

The one-time Perlmutter acceptance check is done; the job scripts now run
on Perlmutter themselves (2d5017e). Its findings (CUDA GPUDirect passes,
VNI order, default CXI service for single-node steps, colocate needs
--overlap, NCCL needs --gpus-per-node=4) move to docs/results.md; the
README job-scripts section explains how to submit on Perlmutter.
Recover the script with: git show 407935b:examples/scripts/perlmutter-check.sh

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
Formatting only (black's AST safety check passed); 8 files.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
All steps are done (cb1761b, 46e006e, 3e5af49); get_batch is documented in
the README and the measurements are in docs/results.md. The plan remains in
history (ddcadd3).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
The v2.0 tag (2025-09-18) was never released; this branch is the 2.0
release instead of 3.0. The tag is to be moved to the release commit.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
…taLoader

pyddstore becomes a package: the extension is pyddstore._core (loaded on
first use, so `from pyddstore.torch import ...` imports torch before MPI
starts); `import pyddstore as dds; dds.PyDDStore` is unchanged. torch is an
optional extra.

pyddstore.torch (moved from examples/vae, generalized):
- DistDataset(source, name, comm, ddstore_width, device, add_device,
  method, handshake_dir): any map-style dataset whose samples are a
  tensor / numpy array / number, a tuple or list of them, or a dict of
  them, with fixed shape and dtype per field. One DDStore variable per
  field (name/key); each rank loads its share. Samples come back with the
  source's structure, types, shapes and dtypes. Schema checked on every
  rank and errors raised on all ranks together. __getitems__ reads a batch
  with one get_batch() per field (DDSTORE_BATCH_GET=0: per sample).
- DistDatasetReader(name, handshake_dir, n_core, device): method-2 extra
  member; reads the fields from {name}.meta.json written by the core group.
- ThreadDataLoader: unchanged behavior (incl. today's fixes).

The VAE examples use the library; their private distdataset.py and
ddstore_dataloader.py are removed. test/test_torch.py covers structures,
field kinds, loaders vs plain source, errors, ddstore_width, GPU, and the
method-2 reader.

Frontier (cxi, 2 nodes): test_torch 19/19 on 4 and 16 ranks; existing
suites pass; vae-ddp loss 8.8960 (method 1 w0/w2/GPU, method 0) and
core/extra split 19.9202 (host and --gpudirect), identical to before.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
chunk_size=N loads each rank's share N samples at a time: every field is
allocated with init() (collective), then filled chunk by chunk with
update() (local), so only one chunk is held besides the store. Peak RSS for
a 400 MiB share on 1 rank: 1202 MiB without chunking, 461 MiB with
chunk_size=20. Host storage only (rejected with add_device for tensor
fields). Default None keeps the whole-share path.

Sample checks are now raised on every rank via one helper
(_raise_on_all), before allocation (first sample's schema) and after
loading (all samples), with no collectives in between.

test_torch: chunk sizes 1 and 4 for every source/method, chunked error
cases, chunk_size validation. Frontier: 36/36 on 4 and 16 ranks.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
A structured record (np.void), structured ndarray or np.recarray of any
shape is stored as raw bytes in a uint8 variable, as projects packing
samples into record arrays did by hand; its dtype (incl. padded/aligned and
nested layouts) goes into the schema, so DistDatasetReader rebuilds it from
{name}.meta.json too. Reads return the same kind of object with the same
layout; ds.dtypes reports the structured dtype.

test_torch: np.void, 1-element structured array, recarray (2,), padded
nested record, dict mixing a record and a tensor; per item, batched,
ThreadDataLoader, chunked; method-2 reader with a record field. Frontier:
57/57 on 4 and 16 ranks.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
- Intro: short feature list (batched reads, GPUDirect, pyddstore.torch,
  thread safety / profiler / split mode).
- Prerequisites/Installation: optional PyTorch, `pip install ".[torch]"`,
  in-place builds need PYTHONPATH=$PWD/src, package layout note and
  rebuild advice for 1.x checkouts.
- Quick Start: pointer to pyddstore.torch.DistDataset.
- Environment variables: new "Read by pyddstore.torch" table
  (DDSTORE_METHOD, DDSTORE_BATCH_GET, handshake dir/timeout, DDSTORE_N_CORE,
  DDSTORE_AFFINITY_*); the examples table keeps only example-only vars.
- Testing: command and table row for test_torch.py (with test_get_batch).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
Raise ValueError ("pass n_core= or set DDSTORE_N_CORE") instead of a bare
KeyError when neither is given.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z
With method 1, epoch_end() and free() don't synchronize, so a rank that
finished its reads could close its endpoint while its peer was still
reading from it. The peer's get() then failed with PTLTE_NOT_FOUND,
skipped all_passed(), and the other rank hung in allreduce. Seen on
Perlmutter in about 1 of 5 runs; 15/15 pass with the barriers.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EasT2JuGMcWVMaVcZXAYxi
- get(): compute the byte count in size_t. disp * itemsize was an int
  product, so rows of 2 GiB or more got a wrapped length and fi_read
  failed with EMSGSIZE.
- Split rows longer than DDSTORE_MAX_READ_BYTES (default 1 GiB, lowered
  to the endpoint's max_msg_size when reported) into several fi_reads.
  On Perlmutter's cxi one 5 GB read fails with EMSGSIZE (2.5 GB works)
  and FI_OPT_MAX_MSG_SIZE reports nothing.
- register_recv(name, arr) / unregister_recv(name, arr): register a
  host or GPU destination buffer once; reads into it or any slice skip
  memory registration. Several per variable, never evicted; the store
  keeps a reference until unregister or free(). No-op for method 0.
- Tests: registered pools from two threads (mr_miss unchanged under
  DDSTORE_PROFILE=1), registered GPU buffer, wide rows (split with
  DDSTORE_MAX_READ_BYTES=4096).

Verified on Perlmutter: test_get_batch (8 and 4 ranks, also with 4 KiB
pieces), test_torch, test_gpu_rdma, test_multirank pass; 5 GB rows read
correctly via get() and get_batch() at 17-20 GB/s.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EasT2JuGMcWVMaVcZXAYxi
…ange

- ThreadDataLoader.__iter__ returns a new _ThreadLoaderIter holding the
  epoch state, whose __iter__ returns itself, as torch's DataLoader does.
  list(it) / islice(it, ...) after next(it) continue the epoch instead of
  restarting it, iterators over one loader are independent (shared thread
  pool), and an iterator dropped early cancels its queued batches.
  next(loader) on the loader itself no longer works (nor does it on
  DataLoader).
- setup.py records the NumPy major version that generated _core.cpp
  (src/pyddstore/_core.numpy-version, ignored) and regenerates the file
  when it changes: a NumPy 2 generated file does not compile against
  NumPy 1.x headers. Also warns about src/pyddstore.cpp / .so left over
  from the old layout.

Verified on Perlmutter: test_torch (8 and 4 ranks, new iterator test
included), test_get_batch, and vae-ddp with 2 loader threads (batched and
per-sample, same loss).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EasT2JuGMcWVMaVcZXAYxi
__next__ refilled before taking its batch, so the slot that batch freed
was only reused at the next __next__: during each training step one
fewer than num_workers * prefetch_factor batches were being fetched
(with 2 threads and prefetch_factor=1, one thread idle every step; in
Walrus this alternated 0.3 s and 2.0-2.4 s data waits). Take the batch,
then refill, as torch's DataLoader does; one more batch is held during
the step.

Test: test_thread_loader_keeps_prefetch_full (3 batches fetched while
one is held with 2 workers, prefetch_factor=1; the old order gave 2).
README: ThreadDataLoader iterator and prefetch depth, registration
wording for get_batch() and free(), tests table.

Verified on Perlmutter: test_torch (8 and 4 ranks), test_get_batch, and
vae-ddp with 2 loader threads (batched and per-sample, same loss).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EasT2JuGMcWVMaVcZXAYxi
For samples made of several stored rows (time windows, clips) and for
reading into reused, registered buffers; on DistDataset and
DistDatasetReader.

- ds.read_rows(rows, fields=None, out=None): any list of stored rows of
  the selected fields, one get_batch() per field, as a dict of
  (len(rows), *shape) values keyed like the sample. Collective for
  method 0, like get_batch().
- ds.alloc(n, fields=None) / ds.release(bufs): one buffer per field,
  register_recv()'d once; read_rows(out=) and __getitems__(idx, out=)
  read into them (views, no registration per read).
- WindowedDataset(ds, window, stride=1, dilation=1, starts=None,
  fields=None): sample i = rows s, s+dilation, ... (s = i*stride or
  starts[i]); a batch of windows is one read_rows().
- row_of(concat, source, index): the row of a source's sample in a store
  built over a ConcatDataset (several files in one store).

Tests: test_read_rows_and_out, test_windowed_dataset, test_row_of_concat
(methods 0/1, DataLoader and ThreadDataLoader). README: PyTorch section
and tests table.

Verified on Perlmutter: test_torch (8 ranks with DDSTORE_PROFILE=1, and
4 ranks on 1 node) and test_get_batch pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EasT2JuGMcWVMaVcZXAYxi
…ffers

reuse_buffers=True reads every batch into one of num_workers buffer sets
from dataset.alloc(batch_size), registered once, instead of fresh buffers
registered on every read. A worker reads and collates its batch, then
returns the set, so the collate must copy: requires batch_size and the
default collate_fn, or collate_copies=True for a custom collate that
copies; other setups (batch_size=None, a non-copying collate, a dataset
without alloc()) raise. The pool belongs to the loader: close() (or
deletion) waits for running fetches and unregisters it, so new loaders
(e.g. one per validation) don't accumulate pinned memory. GPU buffers
are safe because get_batch() synchronizes the device before reading.

Tests: test_thread_loader_reuse_buffers (methods 0/1: two epochs equal
the plain source, no registration per read under DDSTORE_PROFILE=1,
close() unregisters, refusals, copying custom collate) and
test_gpu_reuse_buffers. README: ThreadDataLoader entry, tests table.

Verified on Perlmutter: test_torch (8 ranks with DDSTORE_PROFILE=1, 4
ranks on 1 node) and test_get_batch pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EasT2JuGMcWVMaVcZXAYxi
For samples that can't be stored as they are (strings, labels, metadata
objects, data shared by a group of samples):

- encode(sample) runs on every source sample before it is stored and
  returns what to store.
- decode(stored, index) runs on every sample read (ds[i], __getitems__,
  so in loaders too) and rebuilds the full sample, e.g. decoding ids or
  adding per-group data looked up by a stored id or the index. Runs on
  the reading rank. DistDatasetReader takes decode too.
- fields=[...] keeps only those keys / tuple positions, in order;
  shorthand for an encode, not combinable with one.
- read_rows() and WindowedDataset stay row-level (no decode).

Tests: test_encode_decode (methods 0/1, string labels and a per-group
object through ds[i], __getitems__ and ThreadDataLoader; read_rows
returns stored ids), test_fields_selection, test_method2_reader_decode.
README: DistDataset/DistDatasetReader entries, tests table.

Verified on Perlmutter: test_torch (8 ranks with DDSTORE_PROFILE=1, 4
ranks on 1 node) and test_get_batch pass.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EasT2JuGMcWVMaVcZXAYxi
- docs/: the README's content as pages (installation, quick start,
  backends, GPUDirect, PyTorch integration, HPC systems, performance,
  concurrency, testing, citation) plus references for PyDDStore
  (hand-written, now with get_profile()), pyddstore.torch (autodoc) and
  environment variables. The HPC page adds a table of which Slurm
  --network options each layout needs on Perlmutter and Frontier.
  conf.py mocks mpi4py, torch and the compiled core, so the docs build
  needs no compiler, MPI or libfabric.
- .github/workflows/docs.yml: build with -W on every push and pull
  request; deploy to GitHub Pages from main (needs Settings > Pages >
  Source: GitHub Actions).
- README: 580 -> 117 lines (overview, install, quick start, links to the
  docs pages, citation).
- Consistency: Python >= 3.9 (README, pyproject; the code needs
  importlib.metadata and Executor.shutdown(cancel_futures=)); comments
  naming src/pyddstore.pyx now name src/pyddstore/_core.pyx; README/docs
  show ThreadDataLoader's reuse_buffers and collate_copies.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EasT2JuGMcWVMaVcZXAYxi
@jychoi-hpc
jychoi-hpc merged commit b0b2a5e into main Oct 5, 2026
4 checks passed
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.

1 participant