cupynumeric-parallel-data-load
NVIDIA/skills
Laden Sie fragmentierte Datensätze (npy, Parquet, HDF5, Rohbinärdaten) mithilfe manueller Partitionierung und durch den Start von Legate-Aufgaben in verteilte cuPyNumeric-Arrays.
...Alle erweiternParallel gescharte Daten -> cupynumeric-Ladevorgang
Warum es diese Funktion gibt. cupynumeric spiegelt die Array-API von NumPy wider,
einschließlich cupynumeric.load für eine einzelne .npy-Datei. Darüber hinaus
erfolgt das Laden von Dateien in Legate, nicht in cupynumeric:
| Format | Integrierter Loader |
|---|---|
Einzelne .npy-Datei |
cupynumeric.load(path) (NumPy-API-Kompatibilität) |
| HDF5 (einzelne Datei) | legate.io.hdf5.from_file / from_file_batched |
| Geteilte Mehrdateien (beliebiges Format), Parquet/Arrow, rohe Binärdaten, benutzerdefinierte Layouts | Kein integrierter Loader – diese Funktion. |
Diese Funktion zeigt den kanonischen Weg, um die Lücke in der letzten Zeile zu füllen:
Schreiben Sie eine Legate-Python-Aufgabe, die den für das
Format erforderlichen Leser eines Drittanbieters (h5py, pyarrow, np.memmap, …) innerhalb des
Aufgabenkörpers aufruft, und lassen Sie Legate die Lesevorgänge auf GPUs/Knoten verteilen.
Bei Formaten mit integriertem Loader sollten Sie diesen bevorzugen, es sei denn, Sie benötigen einen
benutzerdefinierten Loader im Task-Körper (mmap-basierter Loader, formatspezifischer Decoder,
Sidecar-Metadaten, partielle/aufgeteilte Lesevorgänge).
Standardmuster: manuelle Partitionierung + manueller Task-Start, dimensioniert auf
die Maschine, nicht auf die Dateien. Nur Achse 0 wird fragmentiert; nachfolgende Achsen
werden innerhalb jeder Kachel mitgeführt. Die Zeilenanzahl pro Shard kann sich zwischen den
Dateien unterscheiden (nur dtype und nachfolgende Achsen müssen übereinstimmen); der Start füllt
jeden verfügbaren Prozessor, unabhängig davon, wie viele Dateien vorhanden sind.
.npy dient als Beispiel, da der Header Form und
Datentyp auf der Festplatte enthält, aber das Grundgerüst gilt für jedes Format mit kostengünstigen
Bereichs-/Slice-Lesevorgängen (Raw-Binär, HDF5, Parquet/Arrow – siehe „Andere
Formate“ weiter unten). Referenzimplementierung:
assets/examples/parallel_npy_load.py.
Annahme zum Datenlayout
Bei dieser Funktion geht es ausschließlich um das Laden – es wird davon ausgegangen, dass die Daten bereits
auf einem gemeinsam genutzten Dateisystem in einer vorhersehbaren, indizierbaren Weise angeordnet sind.
Die Erstellung dieser Dateien liegt außerhalb des Anwendungsbereichs (das Beispiel enthält der Einfachheit halber einen Schreib-
Unterbefehl, aber echte Nutzer stellen ihre eigenen bereit).
Das Beispiel geht von einem bestimmten Layout aus:
- Ein Verzeichnis mit Dateien namens
shard_0000.npy,shard_0001.npy, … in einer zusammenhängenden ganzzahligen Folge (mit Nullen aufgefüllt, Breite 4). - Alle Shards haben denselben
Datentypund dieselben nachgestellten Achsen (shape[1:]); Achse 0 (Zeilen pro Shard) kann sich zwischen den Dateien unterscheiden – das Rezept erstellt eine kumulative Zeilen-Offset-Tabelle und liest den überlappenden Ausschnitt jeder Datei aus der Leaf-Task heraus. - Das Verzeichnis ist für jeden Rank sichtbar (gemeinsames Dateisystem für Läufe mit mehreren Knoten).
Die Funktion `discover_layout()` des Beispiels gibt die Ergebnisse aus und bricht mit einer
ausführlichen Fehlermeldung ab, wenn das Layout falsch ist (fehlendes Verzeichnis,
keine Shards, nicht übereinstimmender Datentyp /nachfolgende Achsen oder eine Lücke in der
zusammenhängenden `shard_NNNN.npy `-Sequenz).
Wenn Ihre Daten in einem anderen Layout vorliegen – als binäre Rohdaten mit festem Schritt, als HDF5-Datei mit einem Datensatz pro Shard, als Verzeichnisbaum usw. – ändern sich lediglich das Glob-Muster, der dateispezifische Reader (Schritt 4 unten) und die Metadaten-Ermittlung (Schritt 1 unten). Die Mechanismen zur Partitionierung und zum Start sind layoutunabhängig.
Wann ist diese Funktion zu verwenden?
Siehe die obige Formattabelle für die Routing-Entscheidung (integrierter Loader vs. diese Funktion). Darüber hinaus gibt es zwei weitere Anhaltspunkte dafür, dass diese Funktion die richtige Wahl ist:
- Ersetzen von sequenziellen
`np.concatenate([read(f) for f in files])` durch parallele Lesevorgänge pro GPU. - Demonstration, wie eine benutzerdefinierte Legate-Python-Aufgabe über einen manuellen Start in ein cupynumerisches Ausgabe-Array schreibt.
Beispiele
Die folgenden Pfade sind relativ zum Verzeichnis dieses Skills angegeben (das Skript
befindet sich unter „assets/examples/parallel_npy_load.py“). Passen Sie das Präfix so an,
dass es mit dem Installationsverzeichnis Ihres Skills übereinstimmt (z. B.
skills/cupynumeric-parallel-data-load/assets/..., wenn sich der Skill
im obersten Verzeichnis „skills/“ befindet).
# Ein Knoten, 4 GPUs.
legate --gpus 4 --fbmem 4000 --min-gpu-chunk 1 \
assets/examples/parallel_npy_load.py \
read --shard-dir /shared/scratch/demo
# Mehrere Knoten, 2 Knoten × 4 GPUs (SLURM), gemeinsames Dateisystem unter --shard-dir.
# Die Shards einmal auf Rank 0 generieren, dann `read` in beliebigem Umfang erneut ausführen.
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
Keine Layout-Flags – der Lese-Treiber durchläuft jeden .npy-Header, um
die Zeilenanzahl pro Datei, die abschließende Form und den Datentyp zu ermitteln, und leitet dann
tile_rows aus der verfügbaren Prozessoranzahl ab.
--min-gpu-chunk 1 ist nur erforderlich, wenn die Elementanzahl pro Kachel
unter der von Legate festgelegten Standard-Mindest-Chunk-Größe für GPU-Starts liegt (z. B. die
Standardeinstellungen des Beispiels – die Gesamtzahl der Zeilen verteilt sich auf 4 GPUs mit
~1 Mio. pro Kachel – den Schwellenwert unterschreiten und andernfalls
auf eine einzige GPU zusammengefasst würden). Bei Datensätzen in Produktionsgröße (zig
Millionen Elemente pro Kachel oder mehr) können Sie das Flag weglassen und
Legate die Standardeinstellung verwenden lassen. Eine Anhebung auf einen moderaten Wert (z. B.
--min-gpu-chunk 1024) ist sinnvoll, wenn jede Kachel groß genug ist, sodass
der Overhead pro Aufgabe wichtiger ist als die Zuweisung einer Kachel pro GPU.
Anleitung
Fünf Schritte aus einem .npy -Beispiel; nur Schritt 1 (Parsen des
Format-Headers) und Schritt 4 (der dateispezifische Reader innerhalb des Aufgaben-Körpers)
sind formatspezifisch. Die anderen drei (Ziel zuweisen, partitionieren,
Fence) werden formatübergreifend unverändert wiederverwendet – siehe „Andere Formate“ weiter unten
für die Swap-Punkte.
1. Metadaten aus jedem Shard lesen
Durchsuche das Verzeichnis und werfe einen Blick auf jeden .npy-Header (mmap_mode="r"
liest nur den Header). Der Header enthält die Shard-spezifische Form und den
Datentyp, sodass der Treiber die Gesamtzahl der Zeilen, die abschließende Form und eine
kumulative Zeilen-Offset-Tabelle ermitteln kann, ohne die Daten jemals zu laden:
paths = sorted(SHARD_DIR.glob("shard_*.npy"))
per_file_rows = [] # Zeilen entlang der Achse 0 pro Datei
trailing_shape = None # shape[1:], muss dateiübergreifend übereinstimmen
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}: Abweichung bei Form und Datentyp "
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) # Länge N+1
total_rows = int(cum_rows[-1])
Der obige Codeausschnitt stellt sicher, dass dtype und trailing_shape (d. h.
shape[1:]) über alle Dateien hinweg übereinstimmen. Die Zeilenanzahl pro Shard kann variieren – die
Tabelle „cum-rows“ berücksichtigt dies. Produktionscode sollte außerdem überprüfen, ob
die Namen eine zusammenhängende Sequenz von shard_0000.npy … shard_NNNN.npy bilden
(im Codeausschnitt der Kürze halber weggelassen; siehe discover_layout() im
Beispiel). Die Erkennung stützt sich ausschließlich auf das, was das
Format auf der Festplatte selbst offenlegt (hier der .npy-Header, .shape /
.dtype für HDF5 usw.); etwaige Begleitdateien (Manifest, Inhalts-Hashes) sind ein
separater Überprüfungsschritt darüber hinaus.
2. Erstellen Sie den cupynumeric-Ausgabespeicher aus den Metadaten
Das Gesamtarray erstreckt sich über total_rows entlang der Achse 0; nachfolgende Achsen stammen
unverändert aus trailing_shape. Verwenden Sie cn.empty – die Aufgabe überschreibt
jede Zelle, eine Initialisierung auf Null wäre verschwendet.
import cupynumeric as cn
total_shape = (total_rows,) + trailing_shape
out = cn.empty(total_shape, dtype=dtype)
3. Aufteilung des Speicherbereichs nach Prozessoranzahl
Die Startform wird an die verfügbaren Prozessoren angepasst, nicht an die
Anzahl der Dateien. Wähle `tile_rows = ceil(total_rows / num_processors) ` und
partitioniere die Achse 0 nach dieser Kachelgröße. Nachfolgende Achsen werden nicht partitioniert
(die Kachel erstreckt sich dort über den gesamten Umfang). Die letzte Kachel darf
kurz sein – genau das unterstützt „partition_by_tiling“ –, sodass das
Rezept keine Teilbarkeitsbeschränkung benötigt.
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 # Anpassung an die Anzahl der Kacheln in der Partition
4. Definieren Sie die Blatt-Aufgabe und starten Sie sie manuell
PATHS und CUM_ROWS (die Dateipfade und die Tabelle mit den kumulativen Zeilen-Offsets
aus Schritt 1) sowie TILE_ROWS werden vor dem Start vom Treiber als globale Modulvariablen
gesetzt; die Replikationssteuerung führt den Treiber
auf jedem Rank aus, sodass jeder Worker identische Werte sieht.
Jede Aufgabe erstellt zunächst ihre Consumer-Ansicht (cupy auf der GPU, numpy auf der
CPU/OMP) und liest die tatsächliche Zeilenanzahl des Tiles aus ` view.shape[0]`
– PhysicalStore selbst verfügt über kein .shape-Attribut, daher ist der Zugriff über
die Ansicht erforderlich. Anschließend berechnet sie ihren globalen Zeilenbereich anhand ihrer
Startkoordinate und dieser Zeilenanzahl, teilt cum_rows für die
überlappenden Dateien in zwei Hälften und kopiert jeden überlappenden Dateischnitt in den
entsprechenden Zielschnitt. Registriere CPU-, OMP- und GPU-Varianten, damit
derselbe Start überall unverändert ausgeführt wird; die Verteilung über
ctx.get_variant_kind() wählt den Consumer aus, der dem Speicherort des
OutputStore ansässig ist (cp.from_dlpack(dst) für FBMEM,
np.asarray(dst) für SYSMEM). cupy wird nur innerhalb des GPU-Zweigs
importiert, sodass der Task-Körper auf Rechnern ohne cupy geladen wird.
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] # Kachelindex 0..num_tasks-1
variant = ctx.get_variant_kind()
if variant == VariantCode.GPU:
import cupy as cp # verzögert: nur auf der GPU
view = cp.from_dlpack(dst)
else:
view = np.asarray(dst) # Zero-Copy-Numpy-Ansicht
tile_rows_actual = view.shape[0] # zu kurz auf der letzten Kachel
row_start = t * TILE_ROWS # globaler Start der Achse 0
row_end = row_start + tile_rows_actual
# Ermitteln des halboffenen Bereichs der Dateiindizes, die sich mit [row_start, row_end) überschneiden.
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):
# Schnittmenge von Kachel [row_start, row_end) mit Datei [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-Schreiben
manual_task = runtime.create_manual_task(
load_tile.library,
load_tile.task_id,
(num_tasks,), # Startdomäne == Anzahl der Kacheln
)
manual_task.add_output(partition)
manual_task.execute()
Beide Verbraucher nutzen die nativen Produzenten von PhysicalStore
(__dlpack__ für cupy, __array_interface__ für np.asarray) –
Zero-Copy-Ansichten der lokalen Kachel. Der Bisect-Aufwand beträgt O(log num_shards)
und die innere Schleife wird typischerweise 1–2 Mal durchlaufen (Kacheln überlappen sich
höchstens über ein paar Dateien).
5. Abgrenzen und Verifizieren
get_legate_runtime().issue_execution_fence(block=True)
Harte Einschränkungen
Alle Shards müssen
denselben Datentypund dieselben nachlaufenden Achsen (shape[1:]) aufweisen. Das Rezept stapelt die Shards entlang der Achse 0; die nachlaufenden Achsen des Ziels stammen austrailing_shape, das im Erkennungsschritt auf den Wert der ersten Datei festgelegt wird. Die Zeilenanzahl pro Shard (shape[0]) darf frei variieren – die Tabelle mit kumulativen Offsets berücksichtigt dies. Das Beispiel lehnt jeden Shard, dessenDatentypoder nachfolgende Form vom ersten abweicht, mit einer aussagekräftigen Fehlermeldung ab.Wählen Sie den Consumer aus, der zur Variante passt.
cp.from_dlpacklehnt SYSMEM-residente Speicher ab;np.asarraygibt stillschweigend eine Host-Ansicht eines FBMEM-residenten Speichers zurück, in den Sie tatsächlich nicht schreiben können. Führen Sie eine Verzweigung anhand von `ctx.get_variant_kind()` durch, damit jede Variante ihren eigenen Consumer verwendet – siehe Schritt 4.mmap-Ansichten sind nicht immer C-kontigu — umschließe jeden dateispezifischen Ausschnitt mit
`np.ascontiguousarray(arr[file_lo:file_hi])`, bevor du`.set()` oder den Numpy-In-Place-Schreibvorgang ausführst.Mehrknoten:
SHARD_DIRmuss sich auf einem gemeinsam genutzten Dateisystem befinden. Jeder Worker (auf jedem Rank) öffnet Shards anhand des Pfads; knotenlokale/tmp-Pfadefunktionieren nur bei Ein-Knoten-Demos.
Varianten
Schneller Pfad mit einheitlichen Shards (eine Aufgabe pro Datei)
Wenn jeder Shard bereits dieselbe Form und denselben Datentyp hat und Sie zufällig
genau so viele Prozessoren wie num_shards zur Verfügung haben, entsteht durch die
„cum-rows“- und „bisect“-Mechanismen ein Overhead. Setze `tile_rows` = `shard_shape[0] ` und
`num_tasks` = `num_shards`; die Partition hat dann eine Kachel pro Datei
und jede Aufgabe liest genau eine Datei von Anfang bis Ende (kein Bisect, keine innere
Schleife). Der Schalter auf der Treiberseite ist eine Einzeiler-Anweisung:
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)
Der gleiche „load_tile“-Aufgaben körper funktioniert in beiden Modi – die innere
Schleife wird lediglich genau einmal pro Aufgabe durchlaufen. Es ist kein
separater Aufgabenkörper für den Fast-Path erforderlich.
Überdekomponierung für eine bessere Lastverteilung
Die Standardeinstellung `tile_rows = ceil(total_rows / num_processors) ` ergibt eine
Kachel pro Prozessor. Um um den Faktor K zu überdekomponieren (kleinere Kacheln,
mehr Punkt-Aufgaben, feinkörnigere Warteschlangen), dividiere stattdessen durch `K * num_processors`:
tile_rows = max(1, (total_rows + K * num_processors - 1) // (K * num_processors))
num_tasks = ceil(total_rows / tile_rows) ergibt dann ungefähr
K * num_processors. Der gleiche Aufgabentext funktioniert weiterhin – die Bisektionsmethode führt lediglich
zu mehr Aufgaben pro Datei.
Andere Formate
Nur der dateispezifische Reader innerhalb von `load_tile` ändert sich. Die
Aufgabe des Readers: Bei Angabe eines Dateipfads und eines halb-offenen Zeilenbereichs
`[file_lo, file_hi)` entlang der Achse 0 soll ein NumPy-Array der Form
(file_hi - file_lo,) + trailing_shape, das C-kontiguös gemacht werden kann.
Es sind kostengünstige Bereichs-/Slice-Lesevorgänge erforderlich – Formate, die nur
„das Lesen der gesamten Datei“ unterstützen, machen den Fall der teilweisen Überlappung zunichte (ein Kachelbereich, der
nur einen Teil einer Datei abdeckt).
| Format | Reader innerhalb der Leaf-Task |
|---|---|
.npy (Beispiel) |
host = np.ascontiguousarray(np.load(p, mmap_mode="r")[file_lo:file_hi]) |
| Rohbinär (feste Form) | arr = np.memmap(p, dtype=DTYPE, mode="r", shape=(rows_in_file, *trailing_shape)); host = np.ascontiguousarray(arr[file_lo:file_hi]) |
| HDF5 | mit h5py.File(p, „r“) als 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 |
(Informationen zu den integrierten Ein-Aufruf-Ladern pro Format finden Sie in der Tabelle „Warum es diese Funktion gibt“ am Anfang dieser Datei.)
Im Erkennungsschritt (Schritt 1) werden die Metadaten jedes Formats analysiert: .npy /
HDF5 / Parquet enthalten auf der Festplatte die Zeilenanzahl pro Datei sowie den Datentyp.
Bei rohen Binärdateien ist dies nicht der Fall – hier wird der Wert aus der Dateigröße abgeleitet oder in einer separaten Datei gespeichert.
Häufige Fallstricke
cn.asarray(dst) ist in einer Blatt-Aufgabe unzulässig
Innerhalb eines @task- Körpers löst jede cupynumeric-Operation, die die oberste
Laufzeitebene berührt – cn.asarray(store), Slice-Zuweisung cn_dst[s] = host_np –
„create_index_space“ aus dem falschen Kontext aus, und Legion bricht ab:
LEGION API USAGE EXCEPTION: Ungültiger Task-Kontext an den Laufzeitaufruf
create_index_space übergeben
Behebung: Verwenden Sie die DLPack-Kapsel mit einer Drittanbieter-Bibliothek (cupy /
torch / numpy) innerhalb von Leaf-Tasks. cn.asarray ist im Treiber zulässig,
jedoch nicht in Leaf-Tasks. Siehe examples/dlpack/leaf_task_interop.py für
die „Torch“-basierte Problemumgehung.
Ein „assert“ innerhalb einer Aufgabe bricht die Laufzeit ab
Legate behandelt nicht ausgelöste Ausnahmen in einem @task als Vertragsverletzung
und bricht den Vorgang ab, sofern die Aufgabe nicht mit throws_exception() registriert wurde.
Führen Sie vor dem Start eine Plausibilitätsprüfung auf dem Host durch.
Die Startdomäne muss mit der Anzahl der Partitionskacheln übereinstimmen
`create_manual_task(launch_shape=...) ` und ` partition_by_tiling(...)`
sind voneinander unabhängig – die Laufzeit erkennt eine Nichtübereinstimmung nicht. Größere Startdomäne
→ Kacheln außerhalb des Bereichs; kleinere → nicht beschriebene Kacheln. Leiten Sie immer
beide aus denselben (total_rows, tile_rows) über zwei separate Ceil-
Divisionen ab (die direkte Anpassung der Startdomäne an num_processors
würde zu einem übermäßigen Start führen, wenn 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,))
```
Alle Dateien
6 Dateiencupynumeric-parallel-data-load installieren
Laden Sie die Skill-Dateien herunter und entpacken Sie sie in Ihr Verzeichnis „.claude/skills/“.
ZIP herunterladenKlonen Sie das Repository und kopieren Sie die Skill-Dateien in Ihr Projekt.
git clone https://github.com/NVIDIA/skills/tree/main/skills/cupynumeric-parallel-data-load # Copy SKILL.md to your .claude/skills/ directory
Kopieren





Heim
