cupynumeric-parallel-data-load
NVIDIA/skills
Cargar conjuntos de datos fragmentados (npy, Parquet, HDF5, binario sin procesar) en matrices distribuidas de cuPyNumeric mediante partición manual y la ejecución de tareas con Legate.
...Expandir todoDatos fragmentados en paralelo -> carga con cupynumeric
Por qué existe esta habilidad. cupynumeric reproduce la API de matrices de NumPy,
incluido cupynumeric.load para un único archivo .npy. Más allá de eso,
la carga de archivos se realiza en Legate, no en cupynumeric:
| Formato | Cargador integrado |
|---|---|
Un único archivo .npy |
cupynumeric.load(path) (compatibilidad con la API de NumPy) |
| HDF5 (archivo único) | legate.io.hdf5.from_file / from_file_batched |
| Varios archivos fragmentados (cualquier formato), Parquet/Arrow, binario sin procesar, diseños personalizados | No hay cargador integrado: esta habilidad. |
Esta habilidad muestra la forma canónica de cubrir el hueco de la última fila:
escribe una tarea de Legate en Python que llame al lector de terceros que
requiera el formato (h5py, pyarrow, np.memmap, …) dentro del
cuerpo de la tarea, y deja que Legate distribuya las lecturas entre las GPU o los nodos.
Para los formatos con un cargador integrado, utilízalo siempre que sea posible, a menos que necesites un
cuerpo de tarea personalizado (cargador basado en mmap, decodificador específico del formato,
metadatos sidecar, lecturas parciales o fragmentadas).
Patrón canónico: partición manual + lanzamiento manual de la tarea, dimensionada en función
de la máquina, no de los archivos. Solo se fragmenta el eje 0; los ejes restantes
se incluyen dentro de cada mosaico. El recuento de filas por fragmento puede variar entre
archivos (solo deben coincidir el tipo de datos y los ejes restantes); el lanzamiento ocupa
todos los procesadores disponibles, independientemente del número de archivos.
.npy es el ejemplo práctico porque el encabezado contiene la forma y el
tipo de datos en el disco, pero el esqueleto se aplica a cualquier formato con
lecturas de rango/segmento económicas (binario sin procesar, HDF5, Parquet/Arrow; véase «Otros
formatos» más abajo). Implementación de referencia:
assets/examples/parallel_npy_load.py.
Supuesto sobre la estructura de los datos
Esta función se centra exclusivamente en la carga: da por hecho que los datos ya están
dispuestos en un sistema de archivos compartido de forma predecible e indexable.
La creación de esos archivos queda fuera del alcance de esta función (el ejemplo incluye un
subcomando de escriturapor comodidad, pero los usuarios reales deben aportar el suyo propio).
El ejemplo práctico asume una estructura específica:
- Un directorio que contiene archivos denominados
shard_0000.npy,shard_0001.npy, ... en una secuencia contigua de números enteros (anchura de 4, rellenada con ceros). - Todos los fragmentos comparten el mismo
tipo de datos (dtype)y los mismos ejes finales (shape[1:]); el eje 0 (filas por fragmento) puede variar entre los archivos — la receta crea una tabla acumulativa de desplazamiento de filas y lee el segmento superpuesto de cada archivo desde el interior de la tarea hoja. - El directorio es visible para todos los rangos (sistema de archivos compartido para ejecuciones con varios nodos).
La función discover_layout() del ejemplo muestra lo que ha encontrado y devuelve un error grave
con un mensaje descriptivo cuando la estructura es incorrecta (falta el directorio,
no hay fragmentos, el tipo de datos o los ejes finales no coinciden, o hay un hueco en la
secuencia contigua shard_NNNN.npy ).
Si tus datos tienen una estructura diferente —binario sin procesar de paso fijo, un archivo HDF5 con un conjunto de datos por fragmento, un árbol de directorios, etc.—, solo cambian el patrón glob, el lector por archivo (paso 4 más abajo) y la detección de metadatos (paso 1 más abajo). El mecanismo de partición y ejecución es independiente del formato de los datos.
Cuándo utilizarla
Consulta la tabla de formatos anterior para decidir el enrutamiento (cargador integrado frente a esta habilidad). Además, hay dos indicios adicionales de que esta habilidad es la opción adecuada:
- Sustituir el código secuencial
`np.concatenate([read(f) for f in files])`por lecturas paralelas por GPU. - Demostrar cómo una tarea Legate Python definida por el usuario escribe en una matriz de salida cupynumeric mediante un lanzamiento manual.
Ejemplos
Las rutas que se indican a continuación están escritas de forma relativa al directorio de esta habilidad (el script
se incluye en assets/examples/parallel_npy_load.py). Ajusta el prefijo para que
coincida con la ubicación en la que esté instalada tu habilidad (por ejemplo,
skills/cupynumeric-parallel-data-load/assets/... si la habilidad se encuentra
en un directorio de nivel superior llamado skills/ ).
# Un solo nodo, 4 GPU.
legate --gpus 4 --fbmem 4000 --min-gpu-chunk 1 \
assets/examples/parallel_npy_load.py \
read --shard-dir /shared/scratch/demo
# Multinodo, 2 nodos x 4 GPU (slurm), sistema de archivos compartido en --shard-dir.
# Genera los fragmentos una vez en el rango 0 y, a continuación, vuelve a ejecutar `read` a cualquier escala.
legate --launcher srun --nodes 2 --cpus 1 \
assets/examples/parallel_npy_load.py \
write --shard-dir /shared/scratch/demo
legate --launcher srun --nodes 2 --ranks-per-node 4 \
--gpus 4 --fbmem 4000 --min-gpu-chunk 1 \
assets/examples/parallel_npy_load.py \
read --shard-dir /shared/scratch/demo
Sin indicadores de disposición: el controlador de lectura recorre cada encabezado .npy para recuperar
el recuento de filas por archivo, la forma final y el tipo de datos, y a continuación calcula
tile_rows a partir del número de procesadores disponibles.
--min-gpu-chunk 1 solo es necesario cuando el recuento de elementos por mosaico es
inferior al tamaño mínimo predeterminado de Legate para los lanzamientos en GPU (por ejemplo, los
valores predeterminados del ejemplo práctico: el total de filas repartidas entre 4 GPU a
~1 millón por mosaico— caen por debajo del umbral y, de no ser así, se
agruparían en una sola GPU). Para conjuntos de datos de tamaño de producción (decenas de
millones de elementos por mosaico o más), puedes omitir la opción y
dejar que Legate utilice su valor por defecto. Aumentarlo a un valor moderado (p. ej.,
--min-gpu-chunk 1024) es adecuado cuando cada bloque es lo suficientemente grande como para que
la sobrecarga por tarea sea más importante que asignar un bloque a cada GPU.
Instrucciones
Cinco pasos a partir de un ejemplo práctico en formato .npy; solo el paso 1 (análisis del
encabezado del formato) y el paso 4 (el lector por archivo dentro del cuerpo de la tarea)
son específicos del formato. Los otros tres (asignación de destino, partición y
barrera) se reutilizan sin cambios en todos los formatos; consulta «Otros formatos» más abajo
para conocer los puntos de intercambio.
1. Leer los metadatos de cada fragmento
Explora el directorio y echa un vistazo a cada encabezado .npy (mmap_mode="r"
solo lee el encabezado). El encabezado contiene la forma y el
tipo de datos por fragmento, por lo que el controlador puede recuperar el número total de filas, la forma final y una
tabla acumulativa de desplazamientos de filas sin necesidad de cargar los datos:
paths = sorted(SHARD_DIR.glob("shard_*.npy"))
per_file_rows = [] # filas a lo largo del eje 0 por archivo
trailing_shape = None # shape[1:], debe coincidir en todos los archivos
dtype = None
for p in paths:
hdr = np.load(p, mmap_mode="r")
if trailing_shape is None:
trailing_shape = tuple(hdr.shape[1:])
dtype = hdr.dtype
elif tuple(hdr.shape[1:]) != trailing_shape or hdr.dtype != dtype:
raise RuntimeError(
f"{p.name}: discrepancia entre la forma final y el tipo de datos "
f"({hdr.shape[1:]}/{hdr.dtype} frente a {trailing_shape}/{dtype})"
)
per_file_rows.append(int(hdr.shape[0]))
cum_rows = np.cumsum([0] + per_file_rows, dtype=np.int64) # longitud N+1
total_rows = int(cum_rows[-1])
El fragmento de código anterior garantiza que el tipo de datos (dtype) y la forma final (trailing_shape) (es decir,
shape[1:]) coincidan en todos los archivos. El recuento de filas por fragmento puede variar; la
tabla cum_rows se encarga de ello. El código de producción también debe verificar que
los nombres formen una secuencia contigua shard_0000.npy ... shard_NNNN.npy
(omitida en el fragmento por brevedad; véase discover_layout() en el
ejemplo práctico). La detección se basa únicamente en lo que
el propio formato en disco expone (la cabecera .npy en este caso, .shape /
.dtype para HDF5, etc.); cualquier archivo complementario (manifiesto, hash de contenido) supone un
paso de verificación adicional.
2. Crear el almacén de salida cupynumeric a partir de los metadatos
La matriz total abarca total_rows a lo largo del eje 0; los ejes restantes provienen
de trailing_shape sin modificaciones. Utiliza cn.empty: la tarea sobrescribe
todas las celdas, por lo que una inicialización a cero sería innecesaria.
import cupynumeric as cn
total_shape = (total_rows,) + trailing_shape
out = cn.empty(total_shape, dtype=dtype)
3. Dividir el almacenamiento en mosaicos según el número de procesadores
La forma de lanzamiento se adapta al número de procesadores disponibles, no al
número de archivos. Elige tile_rows = ceil(total_rows / num_processors) y
divide el eje 0 según ese tamaño de mosaico. Los ejes restantes no se dividen
(el mosaico abarca toda la extensión en ese caso). Se permite que el último mosaico sea
más corto —eso es precisamente lo que admite partition_by_tiling—, por lo que la
receta no necesita ninguna restricción de divisibilidad.
from legate.core import TaskTarget, get_legate_runtime
from legate.core.data_interface import as_logical_array
runtime = get_legate_runtime()
machine = runtime.get_machine()
num_processors = max(
machine.count(TaskTarget.GPU),
machine.count(TaskTarget.OMP),
machine.count(TaskTarget.CPU),
1,
)
tile_rows = max(1, (total_rows + num_processors - 1) // num_processors)
tile_shape = (tile_rows,) + trailing_shape
partition = as_logical_array(out).data.partition_by_tiling(tile_shape)
num_tasks = (total_rows + tile_rows - 1) // tile_rows # hacer coincidir el número de mosaicos con la partición
4. Definir la tarea hoja y ejecutarla manualmente
PATHS y CUM_ROWS (las rutas de los archivos y la tabla de desplazamiento acumulativo de filas
del paso 1), junto con TILE_ROWS, se rellenan como variables globales del módulo
por el controlador antes de la ejecución; el control de replicación ejecuta el controlador
en cada rango, de modo que todos los trabajadores ven valores idénticos.
Cada tarea crea primero su vista de consumidor (cupy en la GPU, numpy en
la CPU/OMP) y lee el recuento real de filas del mosaico a partir de view.shape[0]
— PhysicalStore en sí mismo no tiene el atributo .shape, por lo que es necesario
acceder a través de la vista. A continuación, calcula su rango global de filas a partir de su
coordenada de inicio y ese recuento de filas, divide cum_rows por la mitad para los
archivos superpuestos y copia cada segmento de archivo superpuesto en el
segmento de destino correspondiente. Se registran las variantes de CPU, OMP y GPU para que
el mismo lanzamiento se ejecute sin cambios en cualquier lugar; el envío mediante
ctx.get_variant_kind() selecciona el consumidor que coincida con el lugar donde el
resideel OutputStore (cp.from_dlpack(dst) para FBMEM,
np.asarray(dst) para SYSMEM). cupy se importa únicamente en la rama de la GPU,
por lo que el cuerpo de la tarea se carga en máquinas sin cupy.
import bisect
from legate.core import TaskContext, VariantCode
from legate.core.task import OutputStore, task
@task(variants=(VariantCode.CPU, VariantCode.OMP, VariantCode.GPU))
def load_tile(ctx: TaskContext, dst: OutputStore) -> None:
t = ctx.task_index[0] # índice de mosaico 0..num_tasks-1
variant = ctx.get_variant_kind()
if variant == VariantCode.GPU:
import cupy as cp # diferido: solo en GPU
view = cp.from_dlpack(dst)
else:
view = np.asarray(dst) # vista numpy sin copia
tile_rows_actual = view.shape[0] # falta una fila en el último mosaico
row_start = t * TILE_ROWS # inicio global del eje 0
row_end = row_start + tile_rows_actual
# Hallar el intervalo semiabierto de índices de archivo que se solapan con [row_start, row_end).
first_file = bisect.bisect_right(CUM_ROWS, row_start) - 1
last_file = bisect.bisect_right(CUM_ROWS, row_end - 1) - 1
for f in range(first_file, last_file + 1):
# Intersección del mosaico [row_start, row_end) con el archivo [cum[f], cum[f+1]).
lo = max(row_start, int(CUM_ROWS[f]))
hi = min(fin_fila, int(CUM_ROWS[f + 1]))
límite_inferior_ficha = lo - int(CUM_ROWS[f])
límite_superior_ficha = hi - int(CUM_ROWS[f])
dst_lo = lo - row_start
dst_hi = hi - row_start
chunk = np.ascontiguousarray(
np.load(PATHS[f], mmap_mode="r")[file_lo:file_hi]
)
if variant == VariantCode.GPU:
view[dst_lo:dst_hi].set(chunk) # cudaMemcpyAsync H2D
else:
view[dst_lo:dst_hi] = chunk # escritura numpy sin copia
manual_task = runtime.create_manual_task(
load_tile.library,
load_tile.task_id,
(num_tasks,), # número de dominios de ejecución == número de mosaicos
)
manual_task.add_output(partition)
manual_task.execute()
Ambos consumidores pasan por los productores nativos de PhysicalStore
(__dlpack__ para cupy, __array_interface__ para np.asarray):
vistas sin copia del mosaico local. El coste de la bisección es O(log num_shards)
y el bucle interno suele iterar entre 1 y 2 veces (los mosaicos se solapan como
máximo en un par de archivos).
5. Barrera y verificación
get_legate_runtime().issue_execution_fence(block=True)
Restricciones estrictas
Todos los fragmentos deben compartir
tipo de datos (dtype)y ejes finales (shape[1:]). La receta apila los fragmentos a lo largo del eje 0; los ejes finales del destino proceden detrailing_shape, que la fase de descubrimiento fija en el valor del primer archivo. El número de filas por fragmento (shape[0]) puede variar libremente: la tabla de desplazamiento acumulativo se encarga de gestionarlas. El ejemplo rechaza cualquier fragmento cuyotipo de datos (dtype)o forma final difiera del primero, mostrando un error descriptivo.Elige el consumidor que coincida con la variante.
cp.from_dlpackrechaza los almacenes residentes en SYSMEM;np.asarraydevuelve silenciosamente una vista del host de un almacén residente en FBMEM en el que, en realidad, no se puede escribir . Realiza la distribución segúnctx.get_variant_kind()para que cada variante utilice su propio consumidor — véase el paso 4.Las vistas mmap no siempre son contiguas en C: envuelve cada segmento por archivo con
np.ascontiguousarray(arr[file_lo:file_hi])antes de.set()o de la escritura in situ de numpy.Multinodo:
SHARD_DIRdebe estar en un sistema de archivos compartido. Cada trabajador (en cada rango) abre los fragmentos por ruta; las rutas/tmplocales del nodo solo funcionan para demostraciones de un solo nodo.
Variantes
Ruta rápida con fragmentos uniformes (una tarea por archivo)
Cuando todos los fragmentos tienen ya la misma forma y tipo de datos (shape, dtype) y se dispone
de num_shards procesadores, el mecanismo de cum-rows / bisect
supone una sobrecarga. Establece tile_rows = shard_shape[0] y
num_tasks = num_shards; la partición tendrá entonces un mosaico por archivo
y cada tarea leerá exactamente un archivo de principio a fin (sin bisección, sin
bucle interno). El cambio en el controlador se realiza en una sola línea:
if all(r == per_file_rows[0] for r in per_file_rows) and num_shards == num_processors:
tile_rows = per_file_rows[0]
else:
tile_rows = max(1, (total_rows + num_processors - 1) // num_processors)
El mismo cuerpo de la tarea `load_tile ` sigue funcionando en ambos modos; el bucle interno
simplemente itera exactamente una vez por tarea. No es necesario
un cuerpo de tarea independiente para la ruta rápida.
Descomponer en exceso para un mejor equilibrio de carga
El valor predeterminado tile_rows = ceil(total_rows / num_processors) asigna un
mosaico por procesador. Para sobredescomponer en un factor K (mosaicos más pequeños,
más tareas puntuales, colas más detalladas), divide por K * num_processors
en su lugar:
tile_rows = max(1, (total_rows + K * num_processors - 1) // (K * num_processors))
num_tasks = ceil(total_rows / tile_rows) se amplía entonces hasta aproximadamente
K * num_processors. El mismo cuerpo de la tarea sigue funcionando: la bisección simplemente da como resultado
un mayor número de tareas por archivo.
Otros formatos
Solo cambia el lector por archivo dentro de `load_tile`. El contrato del lector:
dada una ruta de archivo y un rango de filas semiabierto
[file_lo, file_hi) a lo largo del eje 0, devuelve una matriz numpy de forma
(file_hi - file_lo,) + trailing_shape que pueda hacerse contiguo en C.
Se requieren lecturas de rango/segmento económicas: los formatos que solo admiten
«leer todo el archivo» impiden el caso de solapamiento parcial (un mosaico que
solo cubre parte de un archivo).
| Formato | Lector dentro de la tarea hoja |
|---|---|
.npy (ejemplo práctico) |
host = np.ascontiguousarray(np.load(p, mmap_mode="r")[file_lo:file_hi]) |
| Binario sin procesar (forma fija) | arr = np.memmap(p, dtype=DTYPE, mode="r", shape=(rows_in_file, *trailing_shape)); host = np.ascontiguousarray(arr[file_lo:file_hi]) |
| HDF5 | con h5py.File(p, «r») como f: host = np.ascontiguousarray(f["data"][file_lo:file_hi]) |
| Parquet / Arrow | tbl = pq.read_table(p, columns=..., use_threads=False).slice(file_lo, file_hi - file_lo); host = tbl.to_pandas().values |
(Para conocer los cargadores integrados de una sola llamada por formato, consulta la tabla «Por qué existe esta función» que se encuentra al principio de este archivo).
El paso de descubrimiento (paso 1) analiza los metadatos de cada formato: .npy /
HDF5 / Parquet incluyen en el disco el recuento de filas por archivo y el tipo de datos (dtype).
El formato binario sin procesar no los incluye; se deducen del tamaño del archivo o mediante un archivo auxiliar.
Errores habituales
cn.asarray(dst) no es válido en una tarea hoja
Dentro del cuerpo de una @task, cualquier operación cupynumeric que interfiera con el
tiempo de ejecución de nivel superior — cn.asarray(store), asignación de corte cn_dst[s] = host_np —
desencadena create_index_space desde un contexto incorrecto y Legion se interrumpe:
EXCEPCIÓN DE USO DE LA API DE LEGION: Se ha pasado un contexto de tarea no válido a la llamada de tiempo de ejecución
create_index_space
Solución: utiliza la cápsula DLPack con una biblioteca de terceros (cupy /
torch / numpy) dentro de las tareas hoja. cn.asarray funciona correctamente en el controlador,
pero no en las tareas hoja. Consulta examples/dlpack/leaf_task_interop.py para ver
la solución alternativa basada en torch.
Una afirmación dentro de la tarea interrumpe el tiempo de ejecución
Legate trata las excepciones no lanzadas en una @task como una violación del contrato
y la interrumpe, a menos que la tarea se haya registrado con throws_exception().
Realiza una comprobación de validación en el host antes de iniciarla.
El dominio de lanzamiento debe coincidir con el recuento de mosaicos de partición
create_manual_task(launch_shape=...) y partition_by_tiling(...)
son independientes: el tiempo de ejecución no detecta una discrepancia. Un dominio de lanzamiento
más grande → mosaicos fuera de rango; más pequeño → mosaicos sin escribir. Deriva siempre
ambos a partir del mismo (total_rows, tile_rows) mediante dos divisiones por el techo
separadas (ajustar el tamaño del dominio de lanzamiento directamente a num_processors
provocaría un lanzamiento excesivo cuando num_processors > total_rows):
tile_rows = max(1, (total_rows + num_processors - 1) // num_processors)
num_tasks = (total_rows + tile_rows - 1) // tile_rows
partition = ...partition_by_tiling((tile_rows,) + trailing_shape)
runtime.create_manual_task(load_tile.library, load_tile.task_id, (num_tasks,))
---
name: cupynumeric-parallel-data-load
description: Load sharded datasets (npy, Parquet, HDF5, raw binary) into distributed cuPyNumeric arrays using manual partitioning and Legate task launches.
license: CC-BY-4.0 OR Apache-2.0
---
# Parallel sharded data -> cupynumeric load
**Why this skill exists.** cupynumeric mirrors NumPy's array API,
including `cupynumeric.load` for a single `.npy` file. Beyond that,
file *loading* lives in Legate, not cupynumeric:
| Format | Built-in loader |
|---|---|
| Single `.npy` | `cupynumeric.load(path)` (NumPy-API parity) |
| HDF5 (single file) | `legate.io.hdf5.from_file` / `from_file_batched` |
| Sharded multi-file (any format), Parquet/Arrow, raw binary, custom layouts | **No built-in loader — this skill.** |
This skill shows the canonical way to fill the gap in the last row:
write a Legate Python task that calls the third-party reader the
format needs (`h5py`, `pyarrow`, `np.memmap`, ...) inside the
task body, and let Legate distribute the reads across GPUs / nodes.
For the formats with a built-in loader, prefer it unless you need a
custom in-task body (mmap-based loader, format-specific decoder,
sidecar metadata, partial / sharded reads).
Canonical pattern: **manual partition + manual task launch, sized to
the machine, not the files.** Only axis 0 is sharded; trailing axes
ride along inside each tile. Per-shard row counts may differ across
files (only `dtype` and trailing axes must match); the launch fills
every available processor regardless of how many files there are.
`.npy` is the worked example because the header carries shape and
dtype on disk, but the skeleton applies to any format with cheap
range/slice reads (raw binary, HDF5, Parquet/Arrow — see "Other
formats" below). Reference implementation:
[`assets/examples/parallel_npy_load.py`](assets/examples/parallel_npy_load.py).
## Data layout assumption
This skill is purely about **loading** — it assumes the data is already
laid out on a shared filesystem in some predictable, indexable way.
Producing those files is out of scope (the example ships a `write`
subcommand for convenience, but real users bring their own).
The worked example assumes one specific layout:
- A directory containing files named `shard_0000.npy`, `shard_0001.npy`,
... in a contiguous integer sequence (zero-padded width 4).
- All shards share the same `dtype` and the same trailing axes
(`shape[1:]`); **axis 0 (rows per shard) may differ across files** —
the recipe builds a cumulative row-offset table and reads each
file's overlapping slice from inside the leaf task.
- The directory is visible to every rank (shared filesystem for
multi-node runs).
The example's `discover_layout()` prints what it found and hard-fails
with a descriptive error when the layout is wrong (missing directory,
no shards, mismatched `dtype` / trailing axes, or a hole in the
contiguous `shard_NNNN.npy` sequence).
If your data lives in a different layout — fixed-stride raw binary, an
HDF5 file with one dataset per shard, a directory tree, ... — only the
glob pattern, the per-file reader (step 4 below), and the metadata
discovery (step 1 below) change. The partitioning and launch machinery
is layout-agnostic.
## When to use
See the format table above for the routing decision (built-in loader
vs. this skill). Beyond that, two additional cues that this skill is
the right fit:
- Replacing sequential `np.concatenate([read(f) for f in files])` with
parallel per-GPU reads.
- Demonstrating how a user-defined Legate Python task writes into a
cupynumeric output array via a manual launch.
## Examples
Paths below are written relative to this skill's directory (the script
ships at `assets/examples/parallel_npy_load.py`). Adjust the prefix to
match wherever your skill is installed (e.g.
`skills/cupynumeric-parallel-data-load/assets/...` if the skill lives
under a top-level `skills/` directory).
```bash
# Single-node, 4 GPUs.
legate --gpus 4 --fbmem 4000 --min-gpu-chunk 1 \
assets/examples/parallel_npy_load.py \
read --shard-dir /shared/scratch/demo
```
```bash
# Multi-node, 2 nodes x 4 GPUs (slurm), shared filesystem at --shard-dir.
# Generate the shards once on rank 0, then re-run `read` at any scale.
legate --launcher srun --nodes 2 --cpus 1 \
assets/examples/parallel_npy_load.py \
write --shard-dir /shared/scratch/demo
legate --launcher srun --nodes 2 --ranks-per-node 4 \
--gpus 4 --fbmem 4000 --min-gpu-chunk 1 \
assets/examples/parallel_npy_load.py \
read --shard-dir /shared/scratch/demo
```
No layout flags — the read driver walks every `.npy` header to recover
per-file row counts, the trailing shape, and the dtype, then derives
`tile_rows` from the available processor count.
`--min-gpu-chunk 1` is only needed when the per-tile element count is
below Legate's default minimum chunk size for GPU launches (e.g. the
worked example's defaults — total rows split across 4 GPUs at
`~1M` per tile — fall below the threshold and would otherwise be
folded onto a single GPU). For production-sized datasets (tens of
millions of elements per tile or larger) you can drop the flag and
let Legate use its default. Bumping it to a moderate value (e.g.
`--min-gpu-chunk 1024`) is fine when each tile is large enough that
per-task overhead matters more than getting *every* GPU a tile.
## Instructions
Five steps from a `.npy` worked example; only step 1 (parsing the
format header) and step 4 (the per-file reader inside the task body)
are format-specific. The other three (allocate destination, partition,
fence) are reused unchanged across formats — see "Other formats" below
for the swap-points.
### 1. Read the metadata from every shard
Scan the directory and peek at every `.npy` header (`mmap_mode="r"`
reads only the header). The header carries the per-shard shape and
dtype, so the driver can recover total rows, trailing shape, and a
cumulative row-offset table without ever loading the data:
```python
paths = sorted(SHARD_DIR.glob("shard_*.npy"))
per_file_rows = [] # rows along axis 0 per file
trailing_shape = None # shape[1:], must match across files
dtype = None
for p in paths:
hdr = np.load(p, mmap_mode="r")
if trailing_shape is None:
trailing_shape = tuple(hdr.shape[1:])
dtype = hdr.dtype
elif tuple(hdr.shape[1:]) != trailing_shape or hdr.dtype != dtype:
raise RuntimeError(
f"{p.name}: trailing shape / dtype mismatch "
f"({hdr.shape[1:]}/{hdr.dtype} vs {trailing_shape}/{dtype})"
)
per_file_rows.append(int(hdr.shape[0]))
cum_rows = np.cumsum([0] + per_file_rows, dtype=np.int64) # length N+1
total_rows = int(cum_rows[-1])
```
The snippet above enforces matching `dtype` and `trailing_shape` (i.e.
`shape[1:]`) across files. **Per-shard row counts may differ** — the
cum-rows table handles that. Production code should also verify that
names form a contiguous `shard_0000.npy ... shard_NNNN.npy` sequence
(omitted from the snippet for brevity; see `discover_layout()` in the
worked example). Discovery relies only on what the
on-disk format itself exposes (the `.npy` header here, `.shape` /
`.dtype` for HDF5, etc.); any sidecar (manifest, content hashes) is a
separate verification step on top.
### 2. Create the cupynumeric output store from the metadata
The total array spans `total_rows` along axis 0; trailing axes come
from `trailing_shape` unchanged. Use `cn.empty` — the task overwrites
every cell, zero-init would be wasted.
```python
import cupynumeric as cn
total_shape = (total_rows,) + trailing_shape
out = cn.empty(total_shape, dtype=dtype)
```
### 3. Tile the store by processor count
The launch shape is sized to the **available processors**, not to the
file count. Pick `tile_rows = ceil(total_rows / num_processors)` and
partition axis 0 by that tile size. Trailing axes are not partitioned
(tile spans the full extent there). The last tile is allowed to be
short — that's exactly what `partition_by_tiling` supports — so the
recipe needs no divisibility constraint.
```python
from legate.core import TaskTarget, get_legate_runtime
from legate.core.data_interface import as_logical_array
runtime = get_legate_runtime()
machine = runtime.get_machine()
num_processors = max(
machine.count(TaskTarget.GPU),
machine.count(TaskTarget.OMP),
machine.count(TaskTarget.CPU),
1,
)
tile_rows = max(1, (total_rows + num_processors - 1) // num_processors)
tile_shape = (tile_rows,) + trailing_shape
partition = as_logical_array(out).data.partition_by_tiling(tile_shape)
num_tasks = (total_rows + tile_rows - 1) // tile_rows # match partition tile count
```
### 4. Define the leaf task and launch it manually
`PATHS` and `CUM_ROWS` (the file paths and cumulative row-offset
table from step 1) plus `TILE_ROWS` are populated as module globals
by the driver before launching; control replication runs the driver
on every rank, so every worker sees identical values.
Each task builds its consumer view first (cupy on GPU, numpy on
CPU/OMP) and reads the tile's actual row count from `view.shape[0]`
— `PhysicalStore` itself has no `.shape` attribute, so going through
the view is required. It then computes its global row range from its
launch coordinate and that row count, bisects `cum_rows` for the
overlapping file(s), and copies each overlapping file slice into the
matching destination slice. Register CPU, OMP, and GPU variants so
the same launch runs unchanged anywhere; dispatch on
`ctx.get_variant_kind()` picks the consumer matching where the
`OutputStore` is resident (`cp.from_dlpack(dst)` for FBMEM,
`np.asarray(dst)` for SYSMEM). cupy is imported inside the GPU
branch only, so the task body loads on machines without cupy.
```python
import bisect
from legate.core import TaskContext, VariantCode
from legate.core.task import OutputStore, task
@task(variants=(VariantCode.CPU, VariantCode.OMP, VariantCode.GPU))
def load_tile(ctx: TaskContext, dst: OutputStore) -> None:
t = ctx.task_index[0] # tile index 0..num_tasks-1
variant = ctx.get_variant_kind()
if variant == VariantCode.GPU:
import cupy as cp # lazy: only on GPU
view = cp.from_dlpack(dst)
else:
view = np.asarray(dst) # zero-copy numpy view
tile_rows_actual = view.shape[0] # short on the last tile
row_start = t * TILE_ROWS # global axis-0 start
row_end = row_start + tile_rows_actual
# Find the half-open range of file indices that overlap [row_start, row_end).
first_file = bisect.bisect_right(CUM_ROWS, row_start) - 1
last_file = bisect.bisect_right(CUM_ROWS, row_end - 1) - 1
for f in range(first_file, last_file + 1):
# Intersection of tile [row_start, row_end) with file [cum[f], cum[f+1]).
lo = max(row_start, int(CUM_ROWS[f]))
hi = min(row_end, int(CUM_ROWS[f + 1]))
file_lo = lo - int(CUM_ROWS[f])
file_hi = hi - int(CUM_ROWS[f])
dst_lo = lo - row_start
dst_hi = hi - row_start
chunk = np.ascontiguousarray(
np.load(PATHS[f], mmap_mode="r")[file_lo:file_hi]
)
if variant == VariantCode.GPU:
view[dst_lo:dst_hi].set(chunk) # cudaMemcpyAsync H2D
else:
view[dst_lo:dst_hi] = chunk # zero-copy numpy write
manual_task = runtime.create_manual_task(
load_tile.library,
load_tile.task_id,
(num_tasks,), # launch domain == tile count
)
manual_task.add_output(partition)
manual_task.execute()
```
Both consumers go through `PhysicalStore`'s native producers
(`__dlpack__` for cupy, `__array_interface__` for `np.asarray`) —
zero-copy views of the local tile. Bisect cost is `O(log num_shards)`
and the inner loop typically iterates 1–2 times (tiles overlap at
most a couple of files).
### 5. Fence and verify
```python
get_legate_runtime().issue_execution_fence(block=True)
```
## Hard constraints
1. **All shards must share `dtype` and trailing axes (`shape[1:]`).**
The recipe stacks shards along axis 0; the destination's trailing
axes come from `trailing_shape`, which the discovery step locks to
the value of the first file. Per-shard row counts (`shape[0]`) may
freely differ — the cumulative-offset table handles them. The
example rejects any shard whose `dtype` or trailing shape differs
from the first one with a descriptive error.
2. **Pick the consumer that matches the variant.** `cp.from_dlpack`
rejects SYSMEM-resident stores; `np.asarray` silently returns a
host view of an FBMEM-resident store you can't actually write
through. Dispatch on `ctx.get_variant_kind()` so each variant uses
its own consumer — see step 4.
3. **mmap views aren't always C-contiguous** — wrap each per-file
slice with `np.ascontiguousarray(arr[file_lo:file_hi])` before
`.set()` or the numpy in-place write.
4. **Multi-node: `SHARD_DIR` must be on a shared filesystem.** Every
worker (on every rank) opens shards by path; node-local `/tmp` paths
only work for single-node demos.
## Variants
### Uniform-shard fast path (one task per file)
When every shard already has the same `(shape, dtype)` and you happen
to have `num_shards` processors available, the cum-rows / bisect
machinery is overhead. Set `tile_rows = shard_shape[0]` and
`num_tasks = num_shards`; the partition then has one tile per file
and each task reads exactly one file end-to-end (no bisect, no inner
loop). The driver-side switch is a one-liner:
```python
if all(r == per_file_rows[0] for r in per_file_rows) and num_shards == num_processors:
tile_rows = per_file_rows[0]
else:
tile_rows = max(1, (total_rows + num_processors - 1) // num_processors)
```
The same `load_tile` task body still works in either mode — the inner
loop just happens to iterate exactly once per task. There's no need
for a separate task body for the fast path.
### Over-decompose for better load balancing
The default `tile_rows = ceil(total_rows / num_processors)` gives one
tile per processor. To over-decompose by a factor `K` (smaller tiles,
more point tasks, finer-grained queueing), divide by `K * num_processors`
instead:
```python
tile_rows = max(1, (total_rows + K * num_processors - 1) // (K * num_processors))
```
`num_tasks = ceil(total_rows / tile_rows)` then expands to roughly
`K * num_processors`. The same task body still works — bisect just lands
on more tasks per file.
### Other formats
Only the per-file reader inside `load_tile` changes. The reader's
contract: given a file path and a half-open row range
`[file_lo, file_hi)` along axis 0, return a numpy array of shape
`(file_hi - file_lo,) + trailing_shape` that can be made C-contiguous.
Cheap range/slice reads are required — formats that only support
"read the whole file" defeat the partial-overlap case (a tile that
covers only part of one file).
| Format | Reader inside the leaf task |
|---|---|
| **`.npy`** (worked example) | `host = np.ascontiguousarray(np.load(p, mmap_mode="r")[file_lo:file_hi])` |
| **Raw binary** (fixed-shape) | `arr = np.memmap(p, dtype=DTYPE, mode="r", shape=(rows_in_file, *trailing_shape)); host = np.ascontiguousarray(arr[file_lo:file_hi])` |
| **HDF5** | `with h5py.File(p, "r") as f: host = np.ascontiguousarray(f["data"][file_lo:file_hi])` |
| **Parquet / Arrow** | `tbl = pq.read_table(p, columns=..., use_threads=False).slice(file_lo, file_hi - file_lo); host = tbl.to_pandas().values` |
(For built-in single-call loaders per format, see the "Why this skill
exists" table at the top of this file.)
The discovery step (step 1) parses each format's metadata: `.npy` /
HDF5 / Parquet all carry per-file row count + dtype on disk.
Raw binary doesn't — sidecar or derive from file size.
## Common pitfalls
### `cn.asarray(dst)` is illegal in a leaf task
Inside a `@task` body, any cupynumeric op that touches the top-level
runtime — `cn.asarray(store)`, slice assignment `cn_dst[s] = host_np` —
triggers `create_index_space` from the wrong context and Legion aborts:
```
LEGION API USAGE EXCEPTION: Invalid task context passed to runtime call
create_index_space
```
Fix: consume the DLPack capsule with a **third-party** library (cupy /
torch / numpy) inside leaf tasks. `cn.asarray` is fine in the driver,
just not in leaf tasks. See `examples/dlpack/leaf_task_interop.py` for
the torch-flavoured workaround.
### In-task `assert` aborts the runtime
Legate treats unraised exceptions in a `@task` as a contract violation
and aborts unless the task was registered with `throws_exception()`.
Sanity-check on the host before launching.
### Launch domain must match the partition tile count
`create_manual_task(launch_shape=...)` and `partition_by_tiling(...)`
are independent — the runtime doesn't catch a mismatch. Larger launch
domain → out-of-range tiles; smaller → unwritten tiles. Always derive
both from the same `(total_rows, tile_rows)` via two separate `ceil`
divisions (sizing the launch domain to `num_processors` directly
would over-launch when `num_processors > total_rows`):
```python
tile_rows = max(1, (total_rows + num_processors - 1) // num_processors)
num_tasks = (total_rows + tile_rows - 1) // tile_rows
partition = ...partition_by_tiling((tile_rows,) + trailing_shape)
runtime.create_manual_task(load_tile.library, load_tile.task_id, (num_tasks,))
```
Todos los archivos
6 archivosInstalar cupynumeric-parallel-data-load
Descarga y descomprime los archivos de habilidades en tu directorio .claude/skills/.
Descargar ZIPClona el repositorio y copia los archivos de la habilidad a tu proyecto.
git clone https://github.com/NVIDIA/skills/tree/main/skills/cupynumeric-parallel-data-load # Copy SKILL.md to your .claude/skills/ directory
Copiar





Hogar
