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 的加载器、特定格式的解码器、
sidecar 元数据、部分/分片读取),否则请优先使用内置加载器。
标准模式:手动分区 + 手动任务启动,规模应根据
机器而非文件来确定。仅轴 0 进行分片;后续轴
在每个分片内保持不变。 各分片中的行数可能因
文件而异(仅数据类型和后续轴必须匹配);任务启动会填满
所有可用处理器,无论文件数量多少。
.npy被用作示例,因为其头文件在磁盘上携带了形状和
数据类型,但该框架适用于任何支持低成本
范围/切片读取的格式(原始二进制、HDF5、Parquet/Arrow — 参见下文“其他
格式”)。 参考实现:
assets/examples/parallel_npy_load.py。
数据布局假设
本技能纯粹涉及加载操作——它假设数据已经
以某种可预测且可索引的方式布局在共享文件系统上。
生成这些文件不在本技能的范围之内(示例中附带了一个写入
子命令以供方便,但实际用户需自行提供)。
示例代码假设了一种特定的布局:
- 一个目录中包含名称为
shard_0000.npy、shard_0001.npy、 ……的文件,文件名按连续的整数序列排列(补零宽度为 4)。 - 所有分片具有相同的
dtype和相同的尾轴 (shape[1:]);轴 0(每个分片的行数)在不同文件之间可能不同—— 该配方会构建一个累积的行偏移量表,并在叶任务内部读取 每个文件的重叠片段。 - 该目录对每个 rank 都是可见的(多节点运行时使用 共享文件系统)。
示例中的discover_layout()会打印检测结果,若布局错误(目录缺失、
无分片、dtype/ 尾轴不匹配,或连续的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 上,
~1M)——若低于该阈值,原本会被
合并到单个 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:])保持一致。各分片中的行数可能不同——
cum-rows 表负责处理这一情况。 生产环境中的代码还应验证
文件名是否构成连续的shard_0000.npy ... shard_NNNN.npy序列
(为简洁起见,代码片段中省略了此部分;请参阅示例中的
discover_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所支持的——因此
该配方无需设置可 divisibility 约束。
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作为模块全局变量
填充;控制复制机制会在每个 rank 上运行驱动程序,
因此每个工作者看到的值完全一致。
每个任务首先构建其消费者视图(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)
硬性约束
所有分片必须具有相同的
dtype和尾部轴(shape[1:])。 该配方沿轴 0 堆叠分片;目标的尾部 轴来自trailing_shape,该参数在发现步骤中被锁定为 第一个文件的值。 各分片行的行数(shape[0])可以 自由不同——累积偏移量表会处理这些差异。该 示例会以描述性错误拒绝任何数据类型或尾部形状与 第一个分片不同的分片。选择与变体匹配的消费者。
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必须位于共享文件系统上。每个 工作进程(在每个 rank 上)都会按路径打开分片;节点本地的/tmp路径 仅适用于单节点演示。
变体
均匀分片快速路径(每个文件一个任务)
当每个分片已具有相同的(形状、数据类型),且恰好
有num_shards个处理器可用时,累积行 / 二分
机制会带来额外开销。 将tile_rows 设为 shard_shape[0],
num_tasks 设为 num_shards;此时该分区中每个文件对应一个分块,
且每个任务会从头到尾读取 exactly one 文件(无需二分法,无需内部
循环)。驱动程序侧的切换只需一行代码:
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,))
```





首页
