選項
首頁首頁 Skill 資料庫管理 cupynumeric-parallel-data-load

cupynumeric-parallel-data-load

NVIDIA/skills NVIDIA/skills

透過手動分區及 Legate 任務啟動,將分片資料集(npy、Parquet、HDF5、原始二進位檔)載入至分散式 cuPyNumeric 陣列中。

...展開全部
0
更新時間 2026-09-28

平行分片資料 → 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)

硬性限制

  1. 所有分片必須具有相同的資料型別和尾部軸(shape[1:])。 此配方會沿著軸 0 堆疊分片;目標的尾部 軸取自trailing_shape,而發現步驟會將其鎖定為 第一個檔案的值。 各分片行的數量(shape[0])可 自由不同——累積偏移量表會處理這些差異。此 範例會以敘述性錯誤訊息,拒絕任何dtype或尾部形狀與 第一個分片不同的分片。

  2. 選擇與變體相符的消費者。 cp.from_dlpack 會拒絕位於 SYSMEM 中的儲存區;而np.asarray則會靜默地返回 位於 FBMEM 中的儲存區的主機視圖,但您實際上無法 透過該視圖進行寫入。 根據ctx.get_variant_kind()進行分派,使每個變體使用 其專屬的消費者 — 請參閱步驟 4。

  3. mmap 檢視不總是 C 連續的— 請在呼叫 .set()或 numpy 就地寫入之前,將每個按檔案劃分的 切片以np.ascontiguousarray(arr[file_lo:file_hi])進行封裝。

  4. 多節點: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,))
在 GitHub 上查看
---
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

複製 複製
快速設定: 將技能資料夾複製到 .claude/skills/ Claude 會自動偵測並使用該技能
儲存庫 NVIDIA/skills

相關技能

microservices-patterns
更新時間 2026-06-29
jpa-patterns
更新時間 2026-06-30
fabric-lakehouse
更新時間 2026-06-30
prisma-expert
更新時間 2026-06-29
OR