Repository navigation
DDStore 3.0: GPUDirect RDMA, batched reads (get_batch), thread-safe get() - #6
Merged
Merged
Conversation
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
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
- 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
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
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.
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-threadintomain, and sets the package version to 3.0 (v2.0is the previous release tag). Suggested merge: squash (the branch history includes earlywip/addcommits); tagv3.0on the merge commit.What's new
GPU-resident buffers (GPUDirect RDMA,
method=1/2,DDSTORE_FABRIC=cxi)add()andget()accept a CUDA/HIPtorch.Tensor: RDMA reads from / writes into GPU memory with no host copy.hsnandmethod=0reject 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'sfi_reads in flight together.method=0: collective, after MDLoader (IPDPSW 2024):MPI_Allgathervof indices +MPI_Alltoallvof rows, in bounded rounds (DDSTORE_ALLTOALL_MAX_BYTES, default 2 MiB).DistDataset/DistDatasetReader.__getitems__use it, so PyTorch'sDataLoaderandThreadDataLoaderread whole batches automatically;DDSTORE_BATCH_GET=0switches back to per-sample reads.Threading
get()is thread-safe: a per-variable lock insideDDStore::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.rcsettings removed: only the main thread calls MPI.Fixes
free()releases the host buffers DDStore allocated (was leakingMPI_Alloc_mem);join()initializes its mutex.job-vae-core-extra.sh(method 2 core/extra as separatesrunsteps) now works on Slingshot:--network=single_node_vni,job_vniplus keeping only the job VNI inSLINGSHOT_VNIS. Defaults to--layout=split-node; colocate works on Perlmutter only (with--overlap).Tools
DDSTORE_PROFILE=1+get_profile(): whereget()/get_batch()time goes.examples/scripts/bench_get.py: per-row latency/throughput vs row size, destination, batch size, threads.--replicate R,--image-scale S;test/test_get_batch.py;examples/scripts/perlmutter-check.sh.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: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
test_single,test_multirank,test_gpu_rdma14/14,test_get_batch20/20 on 4 ranks and 19/19 on 16; VAE losses unchanged across all variants; core/extra split-node (host and--gpudirect).DDSTORE_FABRIC=hsn:test_get_batch19 passed / 1 skipped (cxi-only GPU case),test_gpu_rdma14/14, VAE method 1, GPU buffers refused with a clear error, core/extra split.--overlap).Known limitations
Error configuring interconnect), with or without--overlap; use split-node.#SBATCHlines target Frontier (-A FUS184, 8 ranks × 7 cores per node).🤖 Generated with Claude Code
https://claude.ai/code/session_01RGZFgjHgEW6QidNBZxaE5z