cupynumeric-parallel-data-load
NVIDIA/skills
透過手動分區及 Legate 任務啟動,將分片資料集(npy、Parquet、HDF5、原始二進位檔)載入至分散式 cuPyNumeric 陣列中。
...展開全部平行分片資料 → cupynumeric 載入
此功能存在的理由。cupynumeric 仿照 NumPy 的陣列 API,
包含用於載入單一.npy檔案的cupynumeric.load 函式。除此之外,
檔案載入功能位於 Legate 中,而非 cupynumeric:
| 格式 | 內建載入器 |
|---|---|
單一.npy 檔案 |
cupynumeric.load(path)(與 NumPy API 完全相容) |
| HDF5(單一檔案) | legate.io.hdf5.from_file/from_file_batched |
| 分片式多檔案(任何格式)、Parquet/Arrow、原始二進位檔、自訂佈局 | 無內建載入器 — 請使用此技能。 |
此技能展示了填補最後一行空缺的標準做法:
編寫一個 Legate Python 任務,在
任務本體中呼叫該格式所需的第三方讀取器
(例如h5py、pyarrow、np.memmap 等),並讓 Legate 將讀取作業分發至各 GPU/節點。
對於具有內建載入器的格式,除非您需要
自訂的任務本體(基於 mmap 的載入器、特定格式的解碼器、
附帶元資料、部分/分片讀取),否則請優先使用內建載入器。
標準模式:手動分區 + 手動任務啟動,尺寸依據
機器而非檔案而定。僅軸 0 進行分片;後續軸
則在每個分片內保持不變。 各分片中的行數可能因
檔案而異(僅需dtype和後續軸數一致);啟動時會填滿
所有可用處理器,無論檔案數量為何。
.npy被用作示例,是因為其標頭在磁碟上攜帶形狀與
資料型別,但此框架適用於任何具備低成本
範圍/切片讀取功能的格式(原始二進位、HDF5、Parquet/Arrow — 參見下文「其他
格式」)。 參考實作:
assets/examples/parallel_npy_load.py。
資料佈局假設
此技能純粹涉及載入— 它假設資料已
以某種可預測且可索引的方式佈局於共用檔案系統上。
產生這些檔案超出本範疇(範例為方便起見附帶了一個寫入
子指令,但實際使用者應自行提供)。
此實作範例假設一種特定的佈局:
- 一個目錄中包含名稱依序為
shard_0000.npy、shard_0001.npy、 ... 的檔案,其名稱採用連續的整數序列(以零墊至 4 位寬)。 - 所有分片皆採用相同的
dtype及相同的尾部軸 (shape[1:]);軸 0(每個分片的行數)在不同檔案間可能有所不同— 此食譜會建立一個累積的行偏移量表,並在葉節點任務內部讀取 每個檔案的重疊切片。 - 該目錄對每個 rank 皆可見(多節點執行時採用 共用檔案系統)。
範例中的discover_layout()會輸出偵測到的結果,若佈局錯誤(目錄缺失、
無分片、資料型別/尾軸不符,或連續的shard_NNNN.npy序列中
存在缺口),則會顯示詳盡的錯誤訊息並直接失敗。
若您的資料採用不同的佈局——例如固定步長的原始二進位檔、 每個分片包含一個資料集的 HDF5 檔案、目錄樹等——僅需 調整通配符模式、各檔案的讀取器(下方第 4 步)以及 元資料偵測(下方第 1 步)。 分區與啟動機制 與資料佈局無關。
何時使用
請參閱上方的格式表以決定路由方式(內建載入器 與此技能)。除此之外,還有兩個額外的線索表明此技能 是合適的選擇:
- 將序列化的
np.concatenate([read(f) for f in files])替換為 並行且針對每個 GPU 的讀取。 - 示範如何透過手動啟動,讓使用者自訂的 Legate Python 任務寫入 cupynumeric 輸出陣列。
範例
以下路徑均以本技能的目錄為基準(該腳本
位於assets/examples/parallel_npy_load.py)。 請根據您的技能安裝位置調整前綴
(例如:
skills/cupynumeric-parallel-data-load/assets/...,若技能位於
頂層的skills/目錄下)。
# 單節點,4 顆 GPU。
legate --gpus 4 --fbmem 4000 --min-gpu-chunk 1 \
assets/examples/parallel_npy_load.py \
read --shard-dir /shared/scratch/demo
# 多節點,2 個節點 × 4 張 GPU(slurm),共用檔案系統位於 --shard-dir。
# 先在第 0 個節點生成分片,之後即可在任何規模下重新執行 `read`。
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
無佈局標誌 — 讀取驅動程式會遍歷每個.npy標頭,以復原
每份檔案的列數、尾隨形狀及資料型別,然後根據
可用處理器數量推導出tile_rows。
--min-gpu-chunk 1僅在每塊的元素數
低於 Legate 針對 GPU 啟動設定的預設最小塊大小時才需要(例如
範例中的預設值 — 總行數分攤至 4 個 GPU,每塊約
~100 萬行)會低於閾值,否則將被
合併至單一 GPU 上)。對於生產級規模的資料集(每塊
數千萬個元素或更大),您可以省略此標誌並
讓 Legate 使用其預設值。將其調高至適中數值(例如
--min-gpu-chunk 1024)則在每個區塊足夠大,且
「每項任務的開銷」比「讓每顆GPU 各處理一個區塊」更為重要時,此設定是可行的。
操作說明
來自.npy範例的五個步驟;僅有步驟 1(解析
格式標頭)和步驟 4(任務本體內的每檔讀取器)
是特定於該格式的。 其餘三步驟(分配目標位置、分區、
柵欄)在各種格式間均可直接重複使用 — 請參閱下文「其他格式」
以了解切換點。
1. 從每個分片讀取元資料
掃描目錄並檢視每個.npy標頭(mmap_mode="r"
僅讀取標頭)。標頭包含各分片的形狀與
資料型別,因此驅動程式無需載入資料即可復原總行數、尾部形狀以及
累積行偏移量表:
paths = sorted(SHARD_DIR.glob("shard_*.npy"))
per_file_rows = [] # 每個檔案沿軸 0 的行數
trailing_shape = None # shape[1:],各檔案之間必須一致
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}: 尾部形狀 / 資料型別不匹配 "
f"({hdr.shape[1:]}/{hdr.dtype} 對比 {trailing_shape}/{dtype})"
)
per_file_rows.append(int(hdr.shape[0]))
cum_rows = np.cumsum([0] + per_file_rows, dtype=np.int64) # 長度為 N+1
total_rows = int(cum_rows[-1])
上述程式碼片段確保各檔案之間的dtype和trailing_shape(即
shape[1:])相符。各分片(shard)的行數可能不同——
cum-rows 表格會處理此情況。 生產環境的程式碼還應驗證
名稱是否構成連續的shard_0000.npy ... shard_NNNN.npy序列
(為簡潔起見,此處未於程式碼片段中列出;請參閱範例中的
discovery_layout()函式)。 此發現機制僅依賴於
磁碟格式本身所揭露的資訊(此處為.npy標頭,HDF5 格式則為.shape/
.dtype等);任何輔助檔案(清單檔、內容雜湊值)均屬
其上的獨立驗證步驟。
2. 根據元資料建立 cupynumeric 輸出儲存區
總陣列沿著第 0 軸橫跨total_rows個元素;尾隨軸則
直接取自trailing_shape且未經變更。請使用cn.empty— 此任務會覆寫
每個儲存格,若先進行零初始化將造成資源浪費。
import cupynumeric as cn
total_shape = (total_rows,) + trailing_shape
out = cn.empty(total_shape, dtype=dtype)
3. 根據處理器數量將儲存區進行分區
啟動形狀的大小應依據可用處理器數量而定,而非
檔案數量。選擇tile_rows = ceil(total_rows / num_processors),並
以此分塊大小對軸 0 進行分區。尾部軸不進行分區
(分塊在該處涵蓋完整範圍)。 最後一塊磁磚允許
長度較短 —— 這正是partition_by_tiling所支援的功能 —— 因此
該配方無需設定可除性限制。
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 # 使分區與瓦片數相符
4. 定義葉節點任務並手動啟動
在啟動之前,驅動程式會將PATHS和CUM_ROWS(來自步驟 1 的檔案路徑和累積行偏移
表)加上TILE_ROWS設定為模組的全域變數;複製控制機制會在每個級別上執行驅動程式,
因此每個工作節點所見的值皆相同。
每個任務首先建立其消費者視圖(GPU 上使用 cupy,CPU/OMP 上使用
numpy),並從view.shape[0]讀取該分區的實際行數
— 由於PhysicalStore本身沒有.shape屬性,因此必須
透過視圖獲取。接著,它會根據啟動座標和該行數
計算其全域行範圍,對cum_rows進行二分搜尋以找出
重疊的檔案,並將每個重疊的檔案切片複製到
對應的目的地切片中。 註冊 CPU、OMP 和 GPU 變體,以便
相同的啟動程式在任何地方都能無需修改地運行;透過
ctx.get_variant_kind()進行調度時,會選擇與
OutputStore所在位置(FBMEM 時為cp.from_dlpack(dst),
SYSMEM 時為np.asarray(dst))。cupy 僅在 GPU
分支內部導入,因此任務主體可在未安裝 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] # 磁磚索引 0..num_tasks-1
variant = ctx.get_variant_kind()
if variant == VariantCode.GPU:
import cupy as cp # 懶加載:僅在 GPU 上執行
view = cp.from_dlpack(dst)
else:
view = np.asarray(dst) # 零拷貝 NumPy 檢視
tile_rows_actual = view.shape[0] # 最後一塊磁磚的列數較少
row_start = t * TILE_ROWS # 全局軸 0 的起始位置
row_end = row_start + tile_rows_actual
# 找出與 [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):
# 磁磚 [row_start, row_end) 與檔案 [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 # 零拷貝 numpy 寫入
manual_task = runtime.create_manual_task(
load_tile.library,
load_tile.task_id,
(num_tasks,), # 啟動域數 == 磁磚數
)
manual_task.add_output(partition)
manual_task.execute()
兩個消費者皆透過PhysicalStore 的原生產生器
(cupy 採用__dlpack__ ,np.asarray 採用__array_interface__)—
取得本地區塊的零拷貝檢視。 二分法計算成本為O(log num_shards)
,而內層迴圈通常迭代 1–2 次(磁磚的重疊範圍
至多涵蓋幾個檔案)。
5. 執行柵欄與驗證
get_legate_runtime().issue_execution_fence(block=True)
硬性限制
所有分片必須具有
相同的資料型別和尾部軸(shape[1:])。 此配方會沿著軸 0 堆疊分片;目標的尾部 軸取自trailing_shape,而發現步驟會將其鎖定為 第一個檔案的值。 各分片行的數量(shape[0])可 自由不同——累積偏移量表會處理這些差異。此 範例會以敘述性錯誤訊息,拒絕任何dtype或尾部形狀與 第一個分片不同的分片。選擇與變體相符的消費者。
cp.from_dlpack會拒絕位於 SYSMEM 中的儲存區;而np.asarray則會靜默地返回 位於 FBMEM 中的儲存區的主機視圖,但您實際上無法 透過該視圖進行寫入。 根據ctx.get_variant_kind()進行分派,使每個變體使用 其專屬的消費者 — 請參閱步驟 4。mmap 檢視不總是 C 連續的— 請在呼叫
.set()或 numpy 就地寫入之前,將每個按檔案劃分的 切片以np.ascontiguousarray(arr[file_lo:file_hi])進行封裝。多節點:
SHARD_DIR必須位於共用檔案系統上。每個 工作節點(在每個排名上)皆會根據路徑開啟分片;節點本地的/tmp路徑 僅適用於單節點示範。
變體
均勻分片快速路徑(每個檔案一個任務)
當每個分片已具有相同的(形狀、資料型別),且您剛好
擁有與 num_shards數量相等的處理器時,累加行 / 二分法
機制會造成額外開銷。 將tile_rows 設為 shard_shape[0],並
將 num_tasks 設為 num_shards;此時該分區將針對每個檔案分配一塊瓦片,
且每個任務會從頭到尾精確讀取一個檔案(無二分法,無內層
迴圈)。驅動程式端的切換僅需一行代碼:
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)
相同的load_tile任務本體在兩種模式下皆可運作——內層
迴圈只是恰好在每個任務中迭代一次。無需
為快速路徑另行定義任務本體。
透過過度分解來改善負載平衡
預設的 `tile_rows = ceil(total_rows / num_processors)` 會讓每個處理器
獲得一個瓦片。若要以係數K進行過度分解(瓦片較小、
任務點數較多、排程粒度更細),請改為除以`K * num_processors`:
tile_rows = max(1, (total_rows + K * num_processors - 1) // (K * num_processors))
num_tasks = ceil(total_rows / tile_rows)隨後將擴展至約
K * num_processors。相同的任務本體仍可運作 — 二分法僅會使
每個檔案對應的任務數量增加。
其他格式
僅有load_tile內的每檔讀取器會有所變更。讀取器的
契約如下:給定一個檔案路徑以及沿 0 軸的半開行範圍
[file_lo, file_hi),回傳一個形狀為
(file_hi - file_lo,) + trailing_shape的 NumPy 陣列,且該陣列可被處理為 C 語言的連續記憶體區塊。
必須具備低成本的範圍/切片讀取功能——僅支援
「讀取整個檔案」的格式將無法處理部分重疊的情況(即一個磁磚
僅涵蓋單一檔案的一部分)。
| 格式 | 葉節點任務內的讀取器 |
|---|---|
.npy(實作範例) |
host = np.ascontiguousarray(np.load(p, mmap_mode="r")[file_lo:file_hi]) |
| 原始二進位檔(固定形狀) | arr = np.memmap(p, dtype=DTYPE, mode="r", shape=(rows_in_file, *trailing_shape)); host = np.ascontiguousarray(arr[file_lo:file_hi]) |
| HDF5 | 使用 h5py.File(p, "r") 作為 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 |
(關於各格式內建的單次呼叫載入器,請參閱本檔案頂端的「為何需要此技能 」表格。)
探索階段(步驟 1)會解析各格式的元資料:.npy/
HDF5 / Parquet 皆會在磁碟上儲存每檔的列數 + 資料型別。
原始二進位檔則沒有 — 需透過 sidecar 檔或根據檔案大小推算。
常見陷阱
在葉節點任務中,cn.asarray(dst)是不合法的
在@task本體內,任何觸及頂層
執行環境的 cupynumeric 運算 —cn.asarray(store)、切片賦值cn_dst[s] = host_np—
都會因從錯誤的上下文觸發create_index_space而導致 Legion 中斷:
LEGION API 使用例外:傳遞給執行階段呼叫
create_index_space 的任務上下文無效
解決方法:在葉節點任務中,透過第三方函式庫(cupy /
torch / numpy)來處理 DLPack 膠囊。cn.asarray在驅動程式中
是可行的,但在葉節點任務中則不可行。請參閱examples/dlpack/leaf_task_interop.py以了解
基於 torch 的解決方法。
任務內的assert會導致執行時中斷
Legate 會將@task內未拋出的例外視為合約違規,
並終止執行,除非該任務已透過throws_exception() 進行註冊。
啟動前請在主機端進行合理性檢查。
啟動域必須與分區磁磚數匹配
create_manual_task(launch_shape=...)與partition_by_tiling(...)
是獨立的 — 執行時環境不會攔截不匹配的情況。 啟動
域過大 → 會產生超出範圍的區塊;過小 → 會產生未寫入的區塊。請務必透過兩次獨立的向上取整
除法,從相同的(total_rows, tile_rows)推導出兩者
(若直接將啟動域大小設定為num_processors
,當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,))
```
安裝 cupynumeric-parallel-data-load
請下載並將技能檔案解壓縮至您的 .claude/skills/ 目錄中。
下載 ZIP複製儲存庫並將技能檔案複製到您的專案中。
git clone https://github.com/NVIDIA/skills/tree/main/skills/cupynumeric-parallel-data-load # Copy SKILL.md to your .claude/skills/ directory
複製





首頁
