pyddstore.torch reference

Generated from the docstrings. For how the pieces fit together, see PyTorch integration.

PyTorch integration for DDStore.

  • DistDataset: a map-style torch.utils.data.Dataset backed 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 one PyDDStore.get_batch() per field, which PyTorch’s DataLoader uses automatically.

  • DistDatasetReader: the same, as a method=2 extra member that joins a dataset published by a DistDataset core group through a shared handshake directory.

  • ThreadDataLoader: a DataLoader whose 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 with read_rows(); row_of() maps a sample of one source in a ConcatDataset to 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 k is stored as variable name/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 method 1/2 and DDSTORE_FABRIC=cxi).

  • add_device – keep this rank’s share of tensor fields on this device.

  • method – DDStore backend (default DDSTORE_METHOD or 0).

  • handshake_dir – method=2 directory (default DDSTORE_HANDSHAKE_DIR or ./ddstore_hs).

  • chunk_size – load this rank’s share chunk_size samples 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 (no add_device for tensor fields).

  • encode – encode(sample) -> sample applied 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) -> sample applied 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 by read_rows() or WindowedDataset (row-level reads).

  • fields – store only these keys (dict samples) or positions (tuple/list samples), in this order: shorthand for an encode that selects them. Not together with encode.

ds[i] returns a sample with the source’s structure: tensors stay tensors (on device if given), numpy arrays stay arrays, numbers stay numbers. ds.ddstore is the underlying PyDDStore.

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 to read_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 like alloc() of values shaped (len(rows), *field_shape); scalar fields give 1-D arrays. out: buffers from alloc() (at least len(rows) rows); the values are then views into them.

Rows come back as stored: decode is not applied.

With method=0 this is collective, like get_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 DistDataset published by a method=2 core 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_DIR or ./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’s encode already 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 to read_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 like alloc() of values shaped (len(rows), *field_shape); scalar fields give 1-D arrays. out: buffers from alloc() (at least len(rows) rows); the values are then views into them.

Rows come back as stored: decode is not applied.

With method=0 this is collective, like get_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 DistDataset or DistDatasetReader): time windows, clips, sequences. Each stored row is held once, however many windows use it.

Sample i is rows s, s + dilation, ..., s + (window - 1) * dilation with s = starts[i] if starts is given, else s = 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 one read_rows() (one get_batch() per field), so with method=0 the same collective rule applies as for ds.

pyddstore.torch.row_of(concat, source, index)

The row of sample index of source source in torch.utils.data.ConcatDataset concat, i.e. its index in a DistDataset built 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 DataLoader that 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 as DataLoader; num_workers is 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 most num_workers * prefetch_factor batches are in flight. Random draws match DataLoader’s. DDSTORE_AFFINITY_WIDTH / DDSTORE_AFFINITY_OFFSET pin worker thread i to CPUs [offset + i*width, offset + (i+1)*width) of the process’s affinity.

reuse_buffers=True (dataset with alloc(), e.g. DistDataset): read every batch into one of a fixed pool of num_workers buffer sets from dataset.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: needs batch_size and the default collate_fn, or collate_copies=True to declare that a custom collate_fn copies. 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.