cupynumeric-parallel-data-load
NVIDIA/skills
Charger des ensembles de données fragmentés (npy, Parquet, HDF5, binaire brut) dans des tableaux cuPyNumeric distribués à l'aide d'un partitionnement manuel et du lancement de tâches via Legate.
...Développer toutDonnées fragmentées en parallèle -> chargement via cupynumeric
Pourquoi cette compétence existe-t-elle ? cupynumeric reproduit l'API des tableaux de NumPy,
y compris la fonction cupynumeric.load pour un seul fichier .npy. Au-delà de cela,
le chargement des fichiers s'effectue dans Legate, et non dans cupynumeric :
| Format | Chargeur intégré |
|---|---|
Fichier .npy unique |
cupynumeric.load(path) (parité avec l’API NumPy) |
| HDF5 (fichier unique) | legate.io.hdf5.from_file / from_file_batched |
| Fichiers multiples fragmentés (tout format), Parquet/Arrow, binaire brut, structures personnalisées | Pas de chargeur intégré — cette compétence. |
Cette compétence montre la méthode canonique pour combler la lacune de la dernière ligne :
écrivez une tâche Python Legate qui appelle le lecteur tiers requis par le
format (h5py, pyarrow, np.memmap, ...) à l'intérieur du
corps de la tâche, et laissez Legate répartir les lectures entre les GPU / nœuds.
Pour les formats disposant d’un chargeur intégré, privilégiez-le sauf si vous avez besoin d’un
corps de tâche personnalisé (chargeur basé sur mmap, décodeur spécifique au format,
métadonnées sidecar, lectures partielles / fragmentées).
Modèle canonique : partition manuelle + lancement manuel de la tâche, dimensionné en fonction de
la machine, et non des fichiers. Seul l’axe 0 est fragmenté ; les axes suivants
sont regroupés au sein de chaque tuile. Le nombre de lignes par fragment peut varier d’un
fichier à l’autre (seuls le type de données et les axes suivants doivent correspondre) ; le lancement occupe
tous les processeurs disponibles, quel que soit le nombre de fichiers.
Le format .npy sert d’exemple pratique car l’en-tête contient la forme et le
type de données sur le disque, mais le principe s’applique à tout format permettant des
lectures de plage/tranche peu coûteuses (binaire brut, HDF5, Parquet/Arrow — voir « Autres
formats » ci-dessous). Implémentation de référence :
assets/examples/parallel_npy_load.py.
Hypothèse concernant la disposition des données
Cette compétence concerne uniquement le chargement — elle suppose que les données sont déjà
disposées sur un système de fichiers partagé de manière prévisible et indexable.
La création de ces fichiers n’entre pas dans le champ d’application (l’exemple fournit une sous-commande d’écriture
pour plus de commodité, mais les utilisateurs réels doivent apporter la leur).
L’exemple pratique suppose une structure spécifique :
- un répertoire contenant des fichiers nommés
shard_0000.npy,shard_0001.npy, … selon une séquence d’entiers contiguë (largeur de 4 complétée par des zéros). - Tous les shards partagent le même
dtypeet les mêmes axes de fin (shape[1:]) ; l’axe 0 (nombre de lignes par shard) peut varier d’un fichier à l’autre — la recette construit une table cumulative de décalages de lignes et lit la tranche chevauchante de chaque fichier depuis la tâche feuille. - Le répertoire est accessible à tous les rangs (système de fichiers partagé pour les exécutions multi-nœuds).
La fonction `discover_layout()` de l’exemple affiche ce qu’elle a trouvé et génère une erreur fatale
accompagnée d’un message descriptif lorsque la disposition est incorrecte (répertoire manquant,
absence de fragments, type de données ou axes finaux non correspondants, ou une lacune dans la
séquence contiguë `shard_NNNN.npy `).
Si vos données sont organisées selon une structure différente — format binaire brut à pas fixe, fichier HDF5 avec un ensemble de données par fragment, arborescence de répertoires, etc. — seuls le motif glob, le lecteur par fichier (étape 4 ci-dessous) et la découverte des métadonnées (étape 1 ci-dessous) changent. Le mécanisme de partitionnement et de lancement est indépendant de la structure des données.
Quand l’utiliser
Consultez le tableau des formats ci-dessus pour déterminer le choix du routage (chargeur intégré ou cette compétence). Au-delà de cela, voici deux autres indices indiquant que cette compétence est la solution adaptée :
- Remplacer la séquence
np.concatenate([read(f) for f in files])par des lectures parallèles par GPU. - Démontrer comment une tâche Python Legate définie par l’utilisateur écrit dans un tableau de sortie cupynumeric via un lancement manuel.
Exemples
Les chemins d’accès ci-dessous sont écrits par rapport au répertoire de cette compétence (le script
se trouve dans assets/examples/parallel_npy_load.py). Adaptez le préfixe en
fonction de l’emplacement d’installation de votre skill (par exemple,
skills/cupynumeric-parallel-data-load/assets/… si la skill se trouve
dans un répertoire de niveau supérieur « skills/ » ).
# Un seul nœud, 4 GPU.
legate --gpus 4 --fbmem 4000 --min-gpu-chunk 1 \
assets/examples/parallel_npy_load.py \
read --shard-dir /shared/scratch/demo
# Multi-nœuds, 2 nœuds × 4 GPU (SLURM), système de fichiers partagé à l'emplacement --shard-dir.
# Générer les segments une seule fois sur le nœud 0, puis relancer `read` à n'importe quelle échelle.
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
Aucun indicateur de disposition — le pilote de lecture parcourt chaque en-tête .npy pour récupérer
le nombre de lignes par fichier, la forme finale et le type de données, puis déduit
le nombre de lignes par tuile (tile_rows) à partir du nombre de processeurs disponibles.
L’option --min-gpu-chunk 1 n’est nécessaire que lorsque le nombre d’éléments par tuile est
inférieur à la taille minimale par défaut des chunks définie par Legate pour les lancements sur GPU (par exemple, les
valeurs par défaut de l’exemple pratique — nombre total de lignes réparties sur 4 GPU à
environ 1 million par tuile — tombent en dessous du seuil et seraient sinon
regroupées sur un seul GPU). Pour les jeux de données de taille industrielle (des dizaines de
millions d’éléments par tuile ou plus), vous pouvez omettre ce drapeau et
laisser Legate utiliser sa valeur par défaut. Le régler sur une valeur modérée (par exemple
--min-gpu-chunk 1024) est tout à fait acceptable lorsque chaque tuile est suffisamment grande pour que
la surcharge par tâche prime sur l’attribution d’une tuile à chaque GPU.
Instructions
Cinq étapes tirées d’un exemple pratique au format .npy; seules l’étape 1 (analyse de l’
en-tête de format) et l’étape 4 (le lecteur par fichier à l’intérieur du corps de la tâche)
sont spécifiques au format. Les trois autres (allocation de la destination, partition,
barrière) sont réutilisées telles quelles quel que soit le format — voir « Autres formats » ci-dessous
pour les points de basculement.
1. Lire les métadonnées de chaque fragment
Parcourez le répertoire et examinez chaque en-tête .npy (mmap_mode="r"
ne lit que l’en-tête). L’en-tête contient la forme et le
type de données par fragment, ce qui permet au pilote de récupérer le nombre total de lignes, la forme de fin et un
tableau cumulatif de décalages de lignes sans jamais charger les données :
paths = sorted(SHARD_DIR.glob("shard_*.npy"))
per_file_rows = [] # nombre de lignes le long de l’axe 0 par fichier
trailing_shape = None # shape[1:], doit correspondre d’un fichier à l’autre
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} : incompatibilité entre la forme finale et le type de données "
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) # longueur N+1
total_rows = int(cum_rows[-1])
L'extrait de code ci-dessus garantit la correspondance entre le dtype et le trailing_shape (c'est-à-dire
shape[1:]) d'un fichier à l'autre. Le nombre de lignes par shard peut varier — le
tableau cum_rows gère cette situation. Le code de production doit également vérifier que
les noms forment une séquence contiguë shard_0000.npy … shard_NNNN.npy
(omise dans l’extrait par souci de concision ; voir discover_layout() dans l’
exemple détaillé). La détection repose uniquement sur ce que le
format sur disque lui-même expose (l’en-tête .npy ici, .shape /
.dtype pour HDF5, etc.) ; tout fichier d’accompagnement (manifeste, hachages de contenu) fait l’objet d’une
étape de vérification distincte en plus.
2. Créer le magasin de sortie cupynumeric à partir des métadonnées
Le tableau total s’étend sur total_rows le long de l’axe 0 ; les axes suivants proviennent
de trailing_shape sans modification. Utilisez cn.empty — la tâche écrase
chaque cellule, une initialisation à zéro serait inutile.
import cupynumeric as cn
total_shape = (total_rows,) + trailing_shape
out = cn.empty(total_shape, dtype=dtype)
3. Diviser le magasin en tuiles en fonction du nombre de processeurs
La forme de lancement est dimensionnée en fonction du nombre de processeurs disponibles, et non du
nombre de fichiers. Choisissez tile_rows = ceil(total_rows / num_processors) et
partitionnez l’axe 0 selon cette taille de tuile. Les axes suivants ne sont pas partitionnés
(la tuile s’étend sur toute leur longueur). La dernière tuile peut être
plus courte — c’est exactement ce que prend en charge partition_by_tiling —, de sorte que la
recette ne nécessite aucune contrainte de divisibilité.
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 # faire correspondre le nombre de tuiles à la partition
4. Définir la tâche feuille et la lancer manuellement
PATHS et CUM_ROWS (les chemins d’accès aux fichiers et la table d’offset cumulatif des lignes
de l’étape 1) ainsi que TILE_ROWS sont renseignés en tant que variables globales du module
par le pilote avant le lancement ; la réplication de contrôle exécute le pilote
sur chaque rang, de sorte que chaque travailleur voie des valeurs identiques.
Chaque tâche construit d’abord sa vue consommateur (cupy sur GPU, numpy sur
CPU/OMP) et lit le nombre réel de lignes de la tuile à partir de view.shape[0]
— PhysicalStore ne dispose pas lui-même d’un attribut .shape, il est donc nécessaire de passer par
la vue. Elle calcule ensuite sa plage de lignes globale à partir de sa
coordonnée de lancement et de ce nombre de lignes, divise cum_rows en deux pour les
fichiers qui se chevauchent, et copie chaque tranche de fichier qui se chevauche dans la
tranche de destination correspondante. Enregistrer les variantes CPU, OMP et GPU afin que
le même lancement s’exécute sans modification où que ce soit ; la répartition via
ctx.get_variant_kind() sélectionne le consommateur correspondant à l’endroit où l’
OutputStore réside (cp.from_dlpack(dst) pour FBMEM,
np.asarray(dst) pour SYSMEM). cupy n’est importé que dans la branche GPU,
ce qui permet au corps de la tâche de se charger sur les machines ne disposant pas de 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] # index de tuile 0..num_tasks-1
variant = ctx.get_variant_kind()
if variant == VariantCode.GPU:
import cupy as cp # paresseux : uniquement sur le GPU
view = cp.from_dlpack(dst)
else:
view = np.asarray(dst) # vue NumPy sans copie
tile_rows_actual = view.shape[0] # manque une ligne sur la dernière tuile
row_start = t * TILE_ROWS # début de l'axe 0 global
row_end = row_start + tile_rows_actual
# Trouver l'intervalle semi-ouvert d'indices de fichiers qui chevauche [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 de la tuile [row_start, row_end) avec le fichier [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 # écriture numpy sans copie
manual_task = runtime.create_manual_task(
load_tile.library,
load_tile.task_id,
(num_tasks,), # nombre de domaines de lancement == nombre de tuiles
)
manual_task.add_output(partition)
manual_task.execute()
Les deux consommateurs passent par les producteurs natifs de PhysicalStore
(__dlpack__ pour cupy, __array_interface__ pour np.asarray) —
des vues sans copie de la tuile locale. Le coût de la bisection est de l’ordre de O(log num_shards)
et la boucle interne s’exécute généralement 1 à 2 fois (les tuiles se chevauchent au
maximum sur quelques fichiers).
5. Barrière et vérification
get_legate_runtime().issue_execution_fence(block=True)
Contraintes strictes
Tous les shards doivent partager
le même dtypeet les mêmes axes de fin (shape[1:]). La recette empile les shards le long de l’axe 0 ; les axes de fin de la destination proviennent detrailing_shape, que l’étape de découverte verrouille sur la valeur du premier fichier. Le nombre de lignes par fragment (shape[0]) peut varier librement — la table d’offset cumulatif les gère. L’ exemple rejette tout fragment dontle type de donnéesou la forme des axes finaux diffère de celui du premier, en affichant une erreur descriptive.Choisissez le consommateur qui correspond à la variante.
cp.from_dlpackrejette les magasins résidant en SYSMEM ;np.asarrayrenvoie silencieusement une vue hôte d’un magasin résidant en FBMEM sur lequel vous ne pouvez pas réellement écrire. Effectuez une répartition en fonction dectx.get_variant_kind()afin que chaque variante utilise son propre consommateur — voir l’étape 4.Les vues mmap ne sont pas toujours contiguës en C — enveloppez chaque tranche par fichier avec
`np.ascontiguousarray(arr[file_lo:file_hi])`avant`.set()`ou l’écriture sur place de numpy.Multi-nœuds :
SHARD_DIRdoit se trouver sur un système de fichiers partagé. Chaque travailleur (sur chaque rang) ouvre les fragments par chemin ; les chemins/tmplocaux au nœud ne fonctionnent que pour les démos à nœud unique.
Variantes
Parcours rapide avec des shards uniformes (une tâche par fichier)
Lorsque chaque fragment a déjà la même forme et le même type de données (shape, dtype) et que vous disposez
exactement de num_shards processeurs, le mécanisme cum-rows / bisect
constitue une surcharge. Définissez tile_rows = shard_shape[0] et
num_tasks = num_shards; la partition comporte alors une tuile par fichier
et chaque tâche lit exactement un fichier de bout en bout (pas de bisection, pas de
boucle interne). Le commutateur côté pilote tient en une seule ligne :
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)
Le même corps de tâche `load_tile` fonctionne toujours dans les deux modes — la boucle interne
s'exécute simplement exactement une fois par tâche. Il n’est pas nécessaire
de créer un corps de tâche distinct pour le chemin rapide.
Surdécomposition pour un meilleur équilibrage de charge
La valeur par défaut `tile_rows = ceil(total_rows / num_processors) ` donne une
tuile par processeur. Pour surdécomposer d’un facteur K (tuiles plus petites,
tâches plus nombreuses, mise en file d’attente plus fine), divisez plutôt par `K * num_processors`
:
tile_rows = max(1, (total_rows + K * num_processors - 1) // (K * num_processors))
num_tasks = ceil(total_rows / tile_rows) donne alors environ
K * num_processors. Le même corps de tâche fonctionne toujours — la méthode de bisection aboutit simplement
à un plus grand nombre de tâches par fichier.
Autres formats
Seul le lecteur par fichier à l’intérieur de `load_tile` change. Le contrat du lecteur
est le suivant : étant donné un chemin de fichier et une plage de lignes semi-ouverte
[file_lo, file_hi) le long de l’axe 0, renvoyer un tableau numpy de forme
(file_hi - file_lo,) + trailing_shape pouvant être rendu contigu en C.
Des lectures de plage/tranche peu coûteuses sont requises — les formats qui ne prennent en charge que
la « lecture du fichier entier » empêchent le cas de chevauchement partiel (une tuile qui
ne couvre qu’une partie d’un fichier).
| Format | Lecteur de format à l’intérieur de la tâche feuille |
|---|---|
.npy (exemple concret) |
host = np.ascontiguousarray(np.load(p, mmap_mode="r")[file_lo:file_hi]) |
| Binaire brut (forme fixe) | arr = np.memmap(p, dtype=DTYPE, mode="r", shape=(rows_in_file, *trailing_shape)); host = np.ascontiguousarray(arr[file_lo:file_hi]) |
| HDF5 | avec 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 |
(Pour les chargeurs intégrés à appel unique par format, consultez le tableau « Pourquoi cette fonctionnalité existe » en haut de ce fichier.)
L'étape de découverte (étape 1) analyse les métadonnées de chaque format : .npy /
HDF5 / Parquet contiennent tous le nombre de lignes par fichier et le type de données (dtype) sur le disque.
Ce n'est pas le cas des fichiers binaires bruts — il faut utiliser un fichier sidecar ou déduire ces informations à partir de la taille du fichier.
Pièges courants
cn.asarray(dst) est interdit dans une tâche feuille
À l'intérieur du corps d'une @task, toute opération cupynumeric qui interagit avec l’environnement d’exécution
de niveau supérieur — cn.asarray(store), affectation par tranchon cn_dst[s] = host_np —
déclenche create_index_space depuis un contexte incorrect et Legion s’interrompt :
EXCEPTION D'UTILISATION DE L'API LEGION : contexte de tâche non valide transmis à l'appel du runtime
create_index_space
Solution : utilisez la capsule DLPack avec une bibliothèque tierce (cupy /
torch / numpy) à l’intérieur des tâches feuilles. cn.asarray fonctionne correctement dans le pilote,
mais pas dans les tâches feuilles. Consultez le fichier examples/dlpack/leaf_task_interop.py pour
découvrir la solution de contournement spécifique à torch.
Une assertion au sein d’une tâche provoque l’interruption du runtime
Legate traite les exceptions non levées dans une @task comme une violation de contrat
et interrompt l’exécution, sauf si la tâche a été enregistrée avec throws_exception().
Effectuez un contrôle de cohérence sur l’hôte avant le lancement.
Le domaine de lancement doit correspondre au nombre de tuiles de partition
create_manual_task(launch_shape=...) et partition_by_tiling(...)
sont indépendants — le runtime ne détecte pas de non-correspondance. Domaine de lancement plus grand
→ tuiles hors limites ; plus petit → tuiles non écrites. Déduisez toujours
les deux à partir des mêmes (total_rows, tile_rows) via deux divisions par le plus grand entier
distinctes (dimensionner le domaine de lancement directement en fonction de num_processors
entraînerait un lancement excessif lorsque 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,))
```
Tous les fichiers
6 fichiersInstaller cupynumeric-parallel-data-load
Téléchargez et décompressez les fichiers de compétences dans votre répertoire .claude/skills/.
Télécharger le ZIPClonez le dépôt et copiez les fichiers de compétence dans votre projet.
git clone https://github.com/NVIDIA/skills/tree/main/skills/cupynumeric-parallel-data-load # Copy SKILL.md to your .claude/skills/ directory
Copier





Maison
