pyddstore.torch reference¶
Generated from the docstrings. For how the pieces fit together, see PyTorch integration.
PyTorch integration for DDStore.
DistDataset: a map-styletorch.utils.data.Datasetbacked by DDStore. Each rank loads its share of any map-style source dataset into the store; every rank can then read any sample. Samples keep the source’s structure (a tensor/array/number, a tuple or list of them, or a dict of them) and each field’s shape and dtype.__getitems__reads a whole batch with onePyDDStore.get_batch()per field, which PyTorch’sDataLoaderuses automatically.DistDatasetReader: the same, as amethod=2extra member that joins a dataset published by aDistDatasetcore group through a shared handshake directory.ThreadDataLoader: aDataLoaderwhose workers are threads instead of forked processes (safe with MPI and GPU-resident buffers).WindowedDataset: samples made of several stored rows (time windows, clips), read withread_rows();row_of()maps a sample of one source in aConcatDatasetto its row.
Fields must have the same shape and dtype in every sample (fixed-shape). Supported dtypes: bool, uint8, int32, int64, float32, float64.
Import torch before mpi4py/MPI starts; from pyddstore.torch import ...
does that by itself.
- class pyddstore.torch.DistDataset(source, name, comm=None, ddstore_width=None, device=None, add_device=None, method=None, handshake_dir=None, chunk_size=None, encode=None, decode=None, fields=None)¶
A map-style dataset stored in DDStore across the ranks of comm.
- Parameters:
source – any map-style dataset (
len(source),source[i]). Each rank loads only its contiguous share. Every sample must have the same structure, and each field the same shape and dtype.name – dataset name; field
kis stored as variablename/k.comm – MPI communicator (default
MPI.COMM_WORLD). All its ranks must construct the dataset together.ddstore_width – ranks per independent store (default: all of
comm); each group holds a full copy of the dataset.device – put tensor fields of read samples on this device (GPUDirect RDMA; needs
method1/2 andDDSTORE_FABRIC=cxi).add_device – keep this rank’s share of tensor fields on this device.
method – DDStore backend (default
DDSTORE_METHODor 0).handshake_dir –
method=2directory (defaultDDSTORE_HANDSHAKE_DIRor./ddstore_hs).chunk_size – load this rank’s share
chunk_sizesamples at a time, writing each chunk into the store before reading the next, so only one chunk is held in memory besides the store (default: read the whole share, then add it). Host storage only (noadd_devicefor tensor fields).encode –
encode(sample) -> sampleapplied to every source sample before it is stored: pick and convert what to store (e.g. drop metadata objects, turn a label string into an id). Its result must meet the rules above.decode –
decode(stored, index) -> sampleapplied to every sample read (ds[i],__getitems__), with its index: add back what wasn’t stored (constants, tables looked up by index or by a stored id). Runs on the reading rank, in the loader thread; anything it looks up must be on every rank. Not applied byread_rows()orWindowedDataset(row-level reads).fields – store only these keys (dict samples) or positions (tuple/list samples), in this order: shorthand for an
encodethat selects them. Not together withencode.
ds[i]returns a sample with the source’s structure: tensors stay tensors (ondeviceif given), numpy arrays stay arrays, numbers stay numbers.ds.ddstoreis the underlyingPyDDStore.With
method=0, batched reads are collective: every rank must read the same number of batches (DistributedSampler does that) from one thread.- __getitems__(indices, out=None)¶
A whole batch: one get_batch() per field (DDSTORE_BATCH_GET=0: one get() per sample). Called by DataLoader and ThreadDataLoader. out: buffers from
alloc(); the samples are then views into them.
- alloc(n, fields=None)¶
Buffers for reads of n rows: a dict keyed like the sample (dict key, or tuple/list position, or 0 for a single value) of one buffer per field, each registered once with
register_recv()so reads into it skip memory registration. Pass toread_rows(out=)or__getitems__(out=); reads may use the first rows only. Samples read into these buffers are views into them: the caller decides when a buffer can be reused.release()unregisters them.
- read_rows(rows, fields=None, out=None)¶
Read stored rows rows (global sample indices; any order, repeats allowed) of the selected fields (default: all) with one
get_batch()per field. Returns a dict keyed likealloc()of values shaped(len(rows), *field_shape); scalar fields give 1-D arrays. out: buffers fromalloc()(at leastlen(rows)rows); the values are then views into them.Rows come back as stored:
decodeis not applied.With
method=0this is collective, likeget_batch(): every rank calls it the same number of times, in the same order.
- release(bufs)¶
Unregister buffers from
alloc()(free()does it too).
- class pyddstore.torch.DistDatasetReader(name, handshake_dir=None, n_core=None, device=None, decode=None)¶
A
DistDatasetpublished by amethod=2core group, joined from a separate job (no MPI communicator needed).- Parameters:
name – the core group’s dataset name.
handshake_dir – shared directory (default
DDSTORE_HANDSHAKE_DIRor./ddstore_hs).n_core – number of core ranks (default
DDSTORE_N_CORE).device – put tensor fields of read samples on this device.
decode – as for
DistDataset(the core group’sencodealready ran before storing).
Waits up to
DDSTORE_HANDSHAKE_TIMEOUT_S(default 300 s) for the core group to publish.- __getitems__(indices, out=None)¶
A whole batch: one get_batch() per field (DDSTORE_BATCH_GET=0: one get() per sample). Called by DataLoader and ThreadDataLoader. out: buffers from
alloc(); the samples are then views into them.
- alloc(n, fields=None)¶
Buffers for reads of n rows: a dict keyed like the sample (dict key, or tuple/list position, or 0 for a single value) of one buffer per field, each registered once with
register_recv()so reads into it skip memory registration. Pass toread_rows(out=)or__getitems__(out=); reads may use the first rows only. Samples read into these buffers are views into them: the caller decides when a buffer can be reused.release()unregisters them.
- read_rows(rows, fields=None, out=None)¶
Read stored rows rows (global sample indices; any order, repeats allowed) of the selected fields (default: all) with one
get_batch()per field. Returns a dict keyed likealloc()of values shaped(len(rows), *field_shape); scalar fields give 1-D arrays. out: buffers fromalloc()(at leastlen(rows)rows); the values are then views into them.Rows come back as stored:
decodeis not applied.With
method=0this is collective, likeget_batch(): every rank calls it the same number of times, in the same order.
- release(bufs)¶
Unregister buffers from
alloc()(free()does it too).
- class pyddstore.torch.WindowedDataset(ds, window, stride=1, dilation=1, starts=None, fields=None)¶
Samples made of several stored rows of ds (a
DistDatasetorDistDatasetReader): time windows, clips, sequences. Each stored row is held once, however many windows use it.Sample
iis rowss, s + dilation, ..., s + (window - 1) * dilationwiths = starts[i]if starts is given, elses = i * stride. Use starts to keep only windows that don’t cross a boundary between trajectories or files. Each field comes back stacked, shaped(window, *field_shape), in the structure of ds’s samples, or as a dict of the selected fields.__getitems__reads a whole batch of windows with oneread_rows()(oneget_batch()per field), so withmethod=0the same collective rule applies as fords.
- pyddstore.torch.row_of(concat, source, index)¶
The row of sample index of source source in
torch.utils.data.ConcatDatasetconcat, i.e. its index in aDistDatasetbuilt over concat. Use it to map (file, trajectory, step) to a row when several sources share one store.
- class pyddstore.torch.ThreadDataLoader(dataset, reuse_buffers=False, collate_copies=False, **DataLoader_kwargs)¶
A
DataLoaderthat fetches batches in a thread pool instead of forked worker processes. Threads share the process’s MPI state, CUDA context and Python objects, so it is safe with DDStore and GPU-resident buffers, where forked workers are not. Takes the same arguments asDataLoader;num_workersis the number of threads (0 means 1).Each batch is fetched (via
dataset.__getitems__when present), collated and optionally pinned in a worker thread; at mostnum_workers * prefetch_factorbatches are in flight. Random draws matchDataLoader’s.DDSTORE_AFFINITY_WIDTH/DDSTORE_AFFINITY_OFFSETpin worker thread i to CPUs[offset + i*width, offset + (i+1)*width)of the process’s affinity.reuse_buffers=True(dataset withalloc(), e.g.DistDataset): read every batch into one of a fixed pool ofnum_workersbuffer sets fromdataset.alloc(batch_size), registered once, instead of fresh buffers that are registered on every read. A worker takes a set, reads and collates the batch, and returns the set, so the collate must copy: needsbatch_sizeand the defaultcollate_fn, orcollate_copies=Trueto declare that a customcollate_fncopies.close()(or deleting the loader) waits for running fetches and unregisters the pool.- close()¶
Stop the worker threads; with
reuse_buffers, wait for running fetches first, then unregister the buffer pool. Called on deletion.