heat.dndarray#
Provides Heat’s core data structure, the DNDarray, a distributed n-dimensional array
Module Contents#
- class DNDarray(array: torch.Tensor, gshape: tuple[int, ...], dtype: heat.core.types.datatype, split: int | None, device: heat.core.devices.Device, comm: Communication, balanced: bool)[source]#
Distributed N-Dimensional array. The core element of Heat. It is composed of PyTorch tensors local to each process.
- Parameters:
array (torch.Tensor) – Local array elements
gshape (tuple[int,...]) – The global shape of the array
dtype (datatype) – The datatype of the array
split (int or None) – The axis on which the array is divided between processes
device (Device) – The device on which the local arrays are using (cpu or gpu)
comm (Communication) – The communications object for sending and receiving data
balanced (bool or None) – Describes whether the data are evenly distributed across processes. If this information is not available (
self.balanced is None), it can be gathered via theis_balanced()method (requires communication).
- __array#
- __gshape#
- __dtype#
- __split#
- __device#
- __comm#
- __balanced :bool#
- __ishalo = False#
- __halo_next :torch.Tensor | None = None#
- __halo_prev :torch.Tensor | None = None#
- __partitions_dict__ = None#
- __lshape_map = None#
- __counts_displs = None#
- __prephalo(start, end) torch.Tensor#
Extracts the halo indexed by start, end from
self.arrayin the direction ofself.split- Parameters:
start (int) – Start index of the halo extracted from
self.arrayend (int) – End index of the halo extracted from
self.array
- get_halo(halo_size: int, prev: bool = True, next: bool = True)[source]#
Fetch halos of size
halo_sizefrom neighboring ranks and save them inself.halo_next/self.halo_prev.
- __cat_halo() torch.Tensor#
Return local array concatenated to halos if they are available.
- __array__() numpy.ndarray[source]#
Returns a view of the process-local slice of the
DNDarrayas a numpy ndarray, if theDNDarrayresides on CPU. Otherwise, it returns a copy, on CPU, of the process-local slice ofDNDarrayas numpy ndarray.
- __array_namespace__(*, api_version: str | None = None) Any[source]#
Returns an object that has all the array API functions on it.
- Parameters:
api_version (Optional[str]) – string representing the version of the array API specification to be returned, in
'YYYY.MM'form. If it isNone(default), it returns the namespace corresponding to latest version of the array API specification.
- astype(dtype, copy=True, device: heat.core.devices.Device = None) DNDarray[source]#
Returns a casted version of this array. Casted array is a new array of the same shape but with given type of this array. If copy is
True, the same array is returned instead.- Parameters:
dtype (datatype) – Heat type to which the array is cast
copy (bool, optional) – By default the operation returns a copy of this array. If copy is set to
Falsethe cast is performed in-place and this array is returneddevice (ht.Device, optional) – The device on which to place the array. If
None, keep device. Default: None.
- balance_() None[source]#
Function for balancing a
DNDarraybetween all nodes. To determine if this is needed use theis_balanced()function. If theDNDarrayis already balanced this function will do nothing. This function modifies theDNDarrayitself and will not return anything.Examples
>>> a = ht.zeros((10, 2), split=0) >>> a[:, 0] = ht.arange(10) >>> b = a[3:] [0/2] tensor([[3., 0.], [1/2] tensor([[4., 0.], [5., 0.], [6., 0.]]) [2/2] tensor([[7., 0.], [8., 0.], [9., 0.]]) >>> b.balance_() >>> print(b.gshape, b.lshape) [0/2] (7, 2) (1, 2) [1/2] (7, 2) (3, 2) [2/2] (7, 2) (3, 2) >>> b [0/2] tensor([[3., 0.], [4., 0.], [5., 0.]]) [1/2] tensor([[6., 0.], [7., 0.]]) [2/2] tensor([[8., 0.], [9., 0.]]) >>> print(b.gshape, b.lshape) [0/2] (7, 2) (3, 2) [1/2] (7, 2) (2, 2) [2/2] (7, 2) (2, 2)
- __cast(cast_function) float | int#
Implements a generic cast function for
DNDarrayobjects.- Parameters:
cast_function (function) – The actual cast function, e.g.
floatorint- Raises:
TypeError – If the
DNDarrayobject cannot be converted into a scalar.
- collect_(target_rank: int = 0) None[source]#
A method collecting a distributed DNDarray to one MPI rank, chosen by the target_rank variable. It is a specific case of the
redistribute_method.- Parameters:
target_rank (int, optional) – The rank to which the DNDarray will be collected. Default: 0.
- Raises:
TypeError – If the target rank is not an integer.
ValueError – If the target rank is out of bounds.
Examples
>>> st = ht.ones((50, 81, 67), split=2) >>> print(st.lshape) [0/2] (50, 81, 23) [1/2] (50, 81, 22) [2/2] (50, 81, 22) >>> st.collect_() >>> print(st.lshape) [0/2] (50, 81, 67) [1/2] (50, 81, 0) [2/2] (50, 81, 0) >>> st.collect_(1) >>> print(st.lshape) [0/2] (50, 81, 0) [1/2] (50, 81, 67) [2/2] (50, 81, 0)
- counts_displs() tuple[tuple[int, ...], tuple[int, ...]][source]#
Returns actual counts (number of items per process) and displacements (offsets) of the DNDarray. Does not assume load balance.
- cpu() DNDarray[source]#
Returns a copy of this object in main memory. If this object is already in main memory, then no copy is performed and the original object is returned.
- create_lshape_map(force_check: bool = False) torch.Tensor[source]#
Generate a ‘map’ of the lshapes of the data on all processes. Units are
(process rank, lshape)- Parameters:
force_check (bool, optional) – if False (default) and the lshape map has already been created, use the previous result. Otherwise, create the lshape_map
- create_partition_interface()[source]#
Create a partition interface in line with the DPPY proposal. This is subject to change. The intention of this to facilitate the usage of a general format for the referencing of distributed datasets.
An example of the output and shape is shown below.
- __partitioned__ = {
‘shape’: (27, 3, 2), ‘partition_tiling’: (4, 1, 1), ‘partitions’: {
- (0, 0, 0): {
‘start’: (0, 0, 0), ‘shape’: (7, 3, 2), ‘data’: tensor([…], dtype=torch.int32), ‘location’: [0], ‘dtype’: torch.int32, ‘device’: ‘cpu’
}, (1, 0, 0): {
‘start’: (7, 0, 0), ‘shape’: (7, 3, 2), ‘data’: None, ‘location’: [1], ‘dtype’: torch.int32, ‘device’: ‘cpu’
}, (2, 0, 0): {
‘start’: (14, 0, 0), ‘shape’: (7, 3, 2), ‘data’: None, ‘location’: [2], ‘dtype’: torch.int32, ‘device’: ‘cpu’
}, (3, 0, 0): {
‘start’: (21, 0, 0), ‘shape’: (6, 3, 2), ‘data’: None, ‘location’: [3], ‘dtype’: torch.int32, ‘device’: ‘cpu’
}
}, ‘locals’: [(rank, 0, 0)], ‘get’: lambda x: x,
}
- Return type:
dictionary containing the partition interface as shown above.
- __dlpack__(*args, **kwargs) Any[source]#
Exports the undistributed array for consumption by
from_dlpack()as a DLPack capsule. Any positional arguments*argsand keyword arguments**kwargsare directly forwarded to torch__dlpack__.Note
See Array API for details and the function signature as implemented by torch.
- Raises:
BufferError – if the DNDarray is distributed, as this is not supported by DLPack.
- __dlpack_device__() tuple[enum.Enum, int][source]#
Returns device type and device ID in DLPack format. Meant for use within
from_dlpack().
- fill_diagonal(value: float) DNDarray[source]#
Fill the main diagonal of a 2D
DNDarray. This function modifies the input tensor in-place, and returns the input array.- Parameters:
value (float) – The value to be placed in the
DNDarraysmain diagonal
- __broadcast_value(key: int | tuple[int, ...] | slice, value: DNDarray, **kwargs)#
Broadcasts the assignment DNDarray value to the shape of the indexed array arr[key] if necessary.
- __set(key: int | tuple[int, ...] | list[int], value: float | DNDarray | torch.Tensor)#
Setter for not advanced indexing, i.e. when arr[key] is an in-place view of arr.
- __advanced_setitem_unordered_local(x_local: torch.Tensor, split_key: torch.Tensor, value_torch: torch.Tensor, *, split_axis: int, value_key_start_dim: int, local_offset: int, local_size: int, value_is_scalar: bool, out_dtype: torch.dtype, base_index: tuple | None = None) None#
The function is a helper that updates
x_localin-place according to the logical advanced indexing pattern encoded bysplit_keyand the broadcastedvalue_torch. This helper operates exclusively on localtorch.Tensorviews: -x_localis the local slice of the distributed array on this rank. -split_keycontains GLOBAL indices along the split axis. - Only those indices that fall into[local_offset, local_offset + local_size)are applied on this rank.
- __getitem_scalar(p: ProcessedKey) DNDarray#
Handles single-element extraction. If the scalar index falls on the split axis, the extracted value is broadcasted from the root process to all others.
- __getitem_local(p: ProcessedKey) DNDarray#
Handles process-local indexing (including standard slices and local advanced indices) directly on local array partitions without MPI communication.
- __getitem_descending_slice_distributed(p: ProcessedKey) DNDarray#
Handles negative step slicing along the split axis. This is a workaround as torch does not support negative-step slicing.
- __getitem_mask(p: ProcessedKey) DNDarray#
Handles fast-path boolean masking. Applies the mask locally without requiring MPI communication during extraction, returning a flattened array distributed along the specified split axis.
- __getitem_advanced_distributed(p: ProcessedKey) DNDarray#
Handles advanced indexing with unordered global indices. Defers to
__getitem_unorderedto resolve data dependencies via anAlltoallvexchange.
- __getitem_unordered(key: tuple, output_shape: tuple, output_split: int, out_is_balanced: bool, key_is_mask_like: bool) DNDarray#
Handles the MPI communication (Alltoallv) when the key along the split axis is unordered and indices are global.
- __prepare_unordered_comm(split_key_flat: torch.Tensor, displs: tuple) tuple#
Helper function for distributed unordered indexing. Determines destination ranks, sorts the key, and computes Alltoallv parameters.
- __getitem__(key: Indexer) DNDarray[source]#
Global getter function for DNDarrays.
Returns a new DNDarray corresponding to the selection of values from the original DNDarray as specified by key. The key can be a variety of indexers, including integers, slices, lists, boolean masks, DNDarrays, ndarrays, torch tensors, and a combination thereof.
The function determines the appropriate method to retrieve the requested data based on the type and structure of key, executing MPI communication if the indexing pattern requires data from multiple processes.
Notes
The returned DNDarray will have its shape, split, and balanced status determined according to the indexing operation performed. For more details on supported indexing behaviors, see the indexing documentation.
- Parameters:
key (array-like indexer) – Indices to get from the
DNDarray.
Examples
>>> a = ht.arange(10, split=0) (1/2) >>> tensor([0, 1, 2, 3, 4], dtype=torch.int32) (2/2) >>> tensor([5, 6, 7, 8, 9], dtype=torch.int32) >>> a[1:6] (1/2) >>> tensor([1, 2, 3, 4], dtype=torch.int32) (2/2) >>> tensor([5], dtype=torch.int32) >>> a = ht.zeros((4, 5), split=0) (1/2) >>> tensor([[0., 0., 0., 0., 0.], [0., 0., 0., 0., 0.]]) (2/2) >>> tensor([[0., 0., 0., 0., 0.], [0., 0., 0., 0., 0.]]) >>> a[1:4, 1] (1/2) >>> tensor([0.]) (2/2) >>> tensor([0., 0.])
- gpu() DNDarray#
Returns a copy of this object in GPU memory. If this object is already in GPU memory, then no copy is performed and the original object is returned.
- is_balanced(force_check: bool = False) bool[source]#
Determine if
selfis balanced evenly (or as evenly as possible) across all nodes distributed evenly (or as evenly as possible) across all processes. This is equivalent to returningself.balanced. If no information is available (self.balanced = None), the balanced status will be assessed via collective communication.- Parameters:
force_check (bool, optional) – If True, the balanced status of the
DNDarraywill be assessed via collective communication in any case.
- is_distributed() bool[source]#
Determines whether the data of this
DNDarrayis distributed across multiple processes.
- item()[source]#
Returns the only element of a 1-element
DNDarray. Mirror of the pytorch command by the same name. If size ofDNDarrayis >1 element, then aValueErroris raised (by pytorch)Examples
>>> import heat as ht >>> x = ht.zeros((1)) >>> x.item() 0.0
- numpy() numpy.typing.NDArray[Any][source]#
Returns a copy of the
DNDarrayas numpy ndarray. If theDNDarrayresides on the GPU, the underlying data will be copied to the CPU first.If the
DNDarrayis distributed, an MPI Allgather operation will be performed before converting to np.ndarray, i.e. each MPI process will end up holding a copy of the entire array in memory. Make sure process memory is sufficient!Examples
>>> import heat as ht T1 = ht.random.randn((10,8)) T1.numpy()
- __repr__() str[source]#
Returns a printable representation of the passed DNDarray, targeting developers.
- ravel() DNDarray[source]#
Flattens the
DNDarray.See also
Examples
>>> a = ht.ones((2, 3), split=0) >>> b = a.ravel() >>> a[0, 0] = 4 >>> b DNDarray([4., 1., 1., 1., 1., 1.], dtype=ht.float32, device=cpu:0, split=0)
- redistribute_(lshape_map: torch.Tensor | None = None, target_map: torch.Tensor | None = None) None[source]#
Redistributes the data of the
DNDarrayalong the split axis to match the given target map. This function does not modify the non-split dimensions of theDNDarray. This is an abstraction and extension of the balance function.- Parameters:
lshape_map (torch.Tensor, optional) – The current lshape of processes. Units are
[rank, lshape].target_map (torch.Tensor, optional) – The desired distribution across the processes. Units are
[rank, target lshape]. Note: the only important parts of the target map are the values along the split axis, values which are not along this axis are there to mimic the shape of thelshape_map.
Examples
>>> st = ht.ones((50, 81, 67), split=2) >>> target_map = torch.zeros((st.comm.size, 3), dtype=torch.int64) >>> target_map[0, 2] = 67 >>> print(target_map) [0/2] tensor([[ 0, 0, 67], [0/2] [ 0, 0, 0], [0/2] [ 0, 0, 0]], dtype=torch.int32) [1/2] tensor([[ 0, 0, 67], [1/2] [ 0, 0, 0], [1/2] [ 0, 0, 0]], dtype=torch.int32) [2/2] tensor([[ 0, 0, 67], [2/2] [ 0, 0, 0], [2/2] [ 0, 0, 0]], dtype=torch.int32) >>> print(st.lshape) [0/2] (50, 81, 23) [1/2] (50, 81, 22) [2/2] (50, 81, 22) >>> st.redistribute_(target_map=target_map) >>> print(st.lshape) [0/2] (50, 81, 67) [1/2] (50, 81, 0) [2/2] (50, 81, 0)
- __redistribute_shuffle(snd_pr: int | torch.Tensor, send_amt: int | torch.Tensor, rcv_pr: int | torch.Tensor, snd_dtype: torch.dtype)#
Function to abstract the function used during redistribute for shuffling data between processes along the split axis
- Parameters:
snd_pr (int or torch.Tensor) – Sending process
send_amt (int or torch.Tensor) – Amount of data to be sent by the sending process
rcv_pr (int or torch.Tensor) – Receiving process
snd_dtype (torch.dtype) – Torch type of the data in question
- resplit_(axis: int = None)[source]#
In-place option for resplitting a
DNDarray.- Parameters:
axis (int) – The new split axis,
Nonedenotes gathering, an int will set the new split axis
Examples
>>> a = ht.zeros( ... ( ... 4, ... 5, ... ), ... split=0, ... ) >>> a.lshape (0/2) (2, 5) (1/2) (2, 5) >>> ht.resplit_(a, None) >>> a.split None >>> a.lshape (0/2) (4, 5) (1/2) (4, 5) >>> a = ht.zeros( ... ( ... 4, ... 5, ... ), ... split=0, ... ) >>> a.lshape (0/2) (2, 5) (1/2) (2, 5) >>> ht.resplit_(a, 1) >>> a.split 1 >>> a.lshape (0/2) (4, 3) (1/2) (4, 2)
- __setitem_local(p: ProcessedKey, value: DNDarray, value_is_scalar: bool) None#
Handles process-local item assignment (slices and local indices) directly on local partitions. If value is distributed, MPI communication might be necessary to align it with the target slice before assignment.
- __setitem_descending_slice_distributed(p: ProcessedKey, value: DNDarray, value_is_scalar: bool) None#
Handles assignment via negative-step slicing. Flips the value array and redistributes it to align with the descending split key before performing the local assignment.
- __setitem_mask(p: ProcessedKey, value: DNDarray, value_is_scalar: bool) None#
Handles assignment using boolean masks. If value is distributed, it will be redistributed to match the number of True elements in the local mask before assignment. If value is not distributed, it will be assigned directly to the masked positions on each process, with PyTorch handling any necessary broadcasting.
- __setitem_advanced_distributed(p: ProcessedKey, original_key, value: DNDarray, value_is_scalar: bool, original_split: int = None) None#
Handles advanced indexing assignments where the indexing key is distributed. This method ensures that the value array is properly aligned and redistributed if necessary before performing the local assignment on each process.
- __setitem_unordered(key: tuple | list | torch.Tensor, key_is_mask_like: bool, value: DNDarray, key_is_single_tensor: bool, counts: tuple, displs: tuple, rank: int, key_is_distributed: bool = False) DNDarray#
Handles the MPI communication when assigning a distributed value to a distributed array with unordered global indices.
- __setitem__(key: Indexer, value: float | DNDarray | torch.Tensor)[source]#
Global item setter for DNDarrays.
Assigns values to the specified positions in the
DNDarray. The key can be a variety of indexers, including integers, slices, lists, boolean masks, DNDarrays, ndarrays, torch tensors, or a combination thereof.If a distributed
DNDarrayis given as the value to be set, this function will automatically attempt to align its distribution scheme (split axis and local shapes) with the target indexed array via MPI communication. If the distributions cannot be safely aligned, aValueErrororRuntimeErroris raised.- Parameters:
key (array-like indexer) – Index/indices to be set
value (float | "DNDarray" | torch.Tensor) – Value to be set to the specified positions in the DNDarray (self)
Notes
For more details on supported indexing behaviors, see the indexing documentation.
Examples
>>> a = ht.zeros((4, 5), split=0) (1/2) >>> tensor([[0., 0., 0., 0., 0.], [0., 0., 0., 0., 0.]]) (2/2) >>> tensor([[0., 0., 0., 0., 0.], [0., 0., 0., 0., 0.]]) >>> a[1:4, 1] = 1 >>> a (1/2) >>> tensor([[0., 0., 0., 0., 0.], [0., 1., 0., 0., 0.]]) (2/2) >>> tensor([[0., 1., 0., 0., 0.], [0., 1., 0., 0., 0.]])
- to_device(device: heat.core.devices.Device, /, *, stream: int | Any | None = None) DNDarray[source]#
Copy the array from the device on which it currently resides to the specified
device.- Parameters:
device (Device) – A
Deviceobject.stream (Int or Any, optional) – Stream object to use during copy.
- tolist(keepsplit: bool = False) list[int | float][source]#
Return a copy of the local array data as a (nested) Python list. For scalars, a standard Python number is returned.
- Parameters:
keepsplit (bool) – Whether the list should be returned locally or globally.
Examples
>>> a = ht.array([[0, 1], [2, 3]]) >>> a.tolist() [[0, 1], [2, 3]]
>>> a = ht.array([[0, 1], [2, 3]], split=0) >>> a.tolist() [[0, 1], [2, 3]]
>>> a = ht.array([[0, 1], [2, 3]], split=1) >>> a.tolist(keepsplit=True) (1/2) [[0], [2]] (2/2) [[1], [3]]