Backends¶
MPI RMA (method=0, default)¶
Uses MPI_Win_create and MPI_Get for one-sided remote reads. Works on any MPI-capable cluster without additional hardware. epoch_begin/epoch_end are required to delimit access epochs.
libfabric RDMA (method=1)¶
Uses fi_read for true RDMA transfers over high-speed interconnects (Infiniband/verbs, Cray GNI, Intel PSM2, Cray Slingshot). Lower latency than MPI RMA on supported hardware. epoch_begin/epoch_end are no-ops with this backend.
DDSTORE_FABRIC selects which libfabric provider to open, for method=1/2:
hsn(default, unset) — opens thetcp;ofi_rxmdomain over Cray Slingshot (Frontier).cxi— opens the nativecxidomain over Cray Slingshot (Frontier and Perlmutter; Perlmutter is CXI-only). Required for GPUDirect RDMA, and the default in the Frontier job scripts.
The two are independent code paths (not runtime auto-detection), so set this explicitly per system rather than relying on a guess:
export DDSTORE_FABRIC=hsn # tcp;ofi_rxm (default)
export DDSTORE_FABRIC=cxi # native CXI: Perlmutter, or Frontier with GPUDirect
PyDDStore picks the network interface (FABRIC_IFACE) automatically for method=1/2, based on each rank’s real CPU affinity (os.sched_getaffinity) — no changes needed in your code:
DDSTORE_NIC_MAP— a precomputed CPU→NIC map, used directly if set (no NIC discovery at construction time). Generate it once from a context with reliable NIC visibility, e.g. ansbatchbatch step’s own shell (not a nestedsruntask — NIC/PCI discovery has been observed to fail there), and export it before launching ranks so every one inherits it:export DDSTORE_NIC_MAP=$(python3 -m cpu_nic_map --env) srun ... python train.py
If
DDSTORE_NIC_MAPisn’t set, each rank falls back to a livehwloc-calc/lstopoquery against its own CPU affinity (cpu_nic_map.allocated_nics(), also runnable standalone aspython3 cpu_nic_map.py --allocated) to find the nearest NIC.Set
FABRIC_IFACEexplicitly to override both and force a specific interface, e.g. when the automatic selection picks the wrong one:export FABRIC_IFACE=hsn0 # e.g. Cray Slingshot
Or skip the environment entirely and pass a map straight to the constructor:
PyDDStore(comm, method=1, nic_map="hsn0=0-15,64-79;hsn1=...").
File-based handshake (method=2)¶
Splits the dataset-holding job from the training job entirely: a core group loads and publishes data, and a separate extra group reads it over RDMA (fi_read, same transport as method=1) — the two are independent MPI jobs (e.g. two separate srun/mpirun launches, possibly on different node allocations) that never share a communicator. They rendezvous only through record files written to a shared-filesystem directory (must be visible to all nodes, e.g. Lustre):
Core member — has an MPI communicator, publishes with
add()/init(). Core ranks exchange records viaMPI_Allgather, and rank 0 writes the combined set to a single{name}.binfile (fabric address, MR key, base pointer, row count, dtype per rank) intohandshake_dir.Extra member — no MPI communicator; constructed with
comm_or_none=Noneand an explicitn_core. Callsjoin(name)to poll for and read alln_corecore-rank records, thenget()works exactly as on the core side, reading directly from core-rank memory over RDMA.
# core side — one MPI job
store = dds.PyDDStore(comm, method=2, handshake_dir="/lustre/.../ddstore_hs")
store.add("x", data)
... # wait for the extra side to finish (e.g. a sentinel file)
store.free()
# extra side — a separate MPI job, no comm needed
store = dds.PyDDStore(None, method=2, handshake_dir="/lustre/.../ddstore_hs", n_core=4)
store.join("x")
out = np.zeros((1, ncols), dtype=np.float32)
store.get("x", out, start=global_idx)
store.free()
Environment variables: DDSTORE_HANDSHAKE_DIR, DDSTORE_HANDSHAKE_TIMEOUT_S, DDSTORE_FABRIC and DDSTORE_NIC_MAP — see Environment variables.
See test/test_method2_core.py / test/test_method2_extra.py for a minimal runnable pair, and examples/vae/vae_core_server.py / examples/vae/vae_extra_train.py for a full DDP training example using this split.
ddstore_width grouping (below) is not currently supported with method=2 — every core rank in comm is treated as one group.
Partitioned / Sub-communicator Usage¶
PyDDStore always spans the whole communicator you pass it. To run several independent stores side by side (e.g. one per node), split comm first; each group then holds a full replica of the dataset, partitioned across its own members:
width = 4 # ranks per group, e.g. GPUs per node
sub_comm = comm.Split(rank // width, rank)
store = dds.PyDDStore(sub_comm) # one independent store per group
This keeps sample fetches inside a node at the cost of replicating the data per group. DistDataset exposes it as ddstore_width (None = one store across all ranks). Not supported with method=2.