옵션
집집 Skill 데이터베이스 관리 cupynumeric-parallel-data-load

cupynumeric-parallel-data-load

NVIDIA/skills NVIDIA/skills

수동 파티셔닝 및 Legate 태스크 실행을 통해 샤드화된 데이터 세트(npy, Parquet, HDF5, 원시 바이너리)를 분산 cuPyNumeric 배열로 로드합니다.

...모든 것을 확장하십시오
0
업데이트 된 시간 2026년 9월 28일

병렬 샤드 데이터 -> cupynumeric 로드

이 스킬이 존재하는 이유. cupynumeric은 NumPy의 배열 API를 반영하며, 여기에는 단일 .npy 파일을 위한 cupynumeric.load도 포함됩니다. 그 외에도, 파일 로딩 기능은 cupynumeric이 아닌 Legate에서 처리됩니다:

형식 내장 로더
단일 .npy cupynumeric.load(path) (NumPy API와 호환)
HDF5 (단일 파일) legate.io.hdf5.from_file / from_file_batched
샤드화된 다중 파일(모든 형식), Parquet/Arrow, 원시 바이너리, 사용자 정의 레이아웃 내장 로더 없음 — 이 스킬.

이 스킬은 마지막 행의 공백을 채우는 표준적인 방법을 보여줍니다: 태스크 본문 내에서 해당 형식에 필요한 타사 리더(h5py, pyarrow, np.memmap 등)를 호출하는 Legate Python 태스크를 작성하고, Legate가 GPU/노드 간에 읽기 작업을 분산하도록 합니다. 내장 로더가 있는 형식의 경우, 태스크 본문 내에서 사용자 정의 로더(mmap 기반 로더, 형식별 디코더, 사이드카 메타데이터, 부분/샤딩된 읽기)가 필요하지 않은 한 내장 로더를 우선적으로 사용하십시오.

표준 패턴: 수동 파티셔닝 + 수동 태스크 실행, 파일 크기가 아닌 시스템 사양에 맞춰 조정합니다. 축 0만 샤딩되며, 후행 축들은 각 타일 내에서 함께 처리됩니다. 파일마다 샤드별 행 수가 다를 수 있습니다 ( dtype과 후행 축만 일치하면 됨); 태스크 실행 시 파일 수와 관계없이 사용 가능한 모든 프로세서를 채웁니다.

.npy가 예시로 사용된 이유는 헤더가 디스크에 모양과 dtype을 포함하고 있기 때문이지만, 이 기본 구조는 저렴한 범위/슬라이스 읽기가 가능한 모든 형식(원시 바이너리, HDF5, Parquet/Arrow — 아래 "기타 형식" 참조)에 적용됩니다. 참조 구현: assets/examples/parallel_npy_load.py.

데이터 레이아웃 가정

이 스킬은 순수하게 로딩에만 초점을 맞추며, 데이터가 이미 예측 가능하고 인덱싱 가능한 방식으로 공유 파일 시스템에 배치되어 있다고 가정합니다. 해당 파일을 생성하는 작업은 이 스킬의 범위를 벗어납니다(예제에는 편의를 위해 쓰기하위 명령어가 포함되어 있지만, 실제 사용자는 자체적으로 준비해야 합니다).

실제 예제는 다음과 같은 특정 레이아웃을 가정합니다:

  • shard_0000.npy, shard_0001.npy, ...와 같은 이름의 파일이 연속적인 정수 순서(0으로 채운 너비 4)로 배치된 디렉터리입니다.
  • 모든 샤드는 동일한 dtype과 동일한 후행 축 (shape[1:]) 을 공유합니다. 축 0(샤드당 행 수)은 파일마다 다를 수 있습니다 — 이 레시피는 누적 행 오프셋 테이블을 구축하고, 리프 태스크 내부에서 각 파일의 중첩된 슬라이스를 읽어들입니다.
  • 이 디렉터리는 모든 랭크에서 접근 가능합니다(다중 노드 실행 시 공유 파일 시스템).

이 예제의 discover_layout()은 발견한 내용을 출력하며, 레이아웃이 잘못된 경우(디렉터리가 없거나, 샤드가 없거나, dtype/후행 축이 일치하지 않거나, 연속적인 shard_NNNN.npy 시퀀스에 공백이 있는 경우) 설명이 포함된 오류 메시지와 함께 실행이 중단됩니다.

데이터가 다른 레이아웃(고정 스트라이드 원시 바이너리, 샤드당 하나의 데이터셋을 포함하는 HDF5 파일, 디렉터리 트리 등)에 있는 경우, glob 패턴, 파일별 리더(아래 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/ 디렉터리 아래에 있는 경우).

# 단일 노드, GPU 4개.
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 헤더를 순회하여 파일별 행 수, 후행 셰이프 및 dtype을 복원한 다음, 사용 가능한 프로세서 수를 바탕으로 tile_rows를 도출합니다.

--min-gpu-chunk 1은 타일당 요소 수가 Legate의 GPU 실행에 대한 기본 최소 청크 크기보다 작을 때만 필요합니다(예: 실습 예제의 기본값 — 총 행 수가 4개의 GPU에 분할되어 타일당~1M )가 임계값 아래로 떨어져, 그렇지 않았다면 단일 GPU로 통합되었을 경우)에만 필요합니다. 실제 운영 규모의 데이터셋(타일당 수천만 개 이상의 요소)의 경우 이 플래그를 생략하고 Legate가 기본값을 사용하도록 할 수 있습니다. 이 값을 적당한 수준(예: --min-gpu-chunk 1024)로 설정하는 것은 각 타일이 충분히 커서 모든 GPU에 타일을 할당하는 것보다 작업별 오버헤드가 더 중요한 경우에 적합합니다.

사용 방법

.npy 예제에서 발췌한 5단계 과정입니다. 1단계(포맷 헤더 파싱)와 4단계(작업 본문 내의 파일별 리더)만 포맷에 따라 다릅니다. 나머지 세 단계(대상 할당, 파티셔닝, 펜스)는 형식 간에 변경 없이 재사용됩니다. 교환 지점에 대해서는 아래의 "기타 형식"을 참조하십시오.

1. 모든 샤드에서 메타데이터 읽기

디렉터리를 스캔하여 모든 .npy 헤더를 엿봅니다(mmap_mode="r" 은 헤더만 읽습니다). 헤더에는 샤드별 모양과 dtype이 포함되어 있으므로, 드라이버는 데이터를 로드하지 않고도 총 행 수, 후행 모양 및 누적 행 오프셋 테이블을 복원할 수 있습니다:

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} vs {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를 사용하세요 — 이 작업은 모든 셀을 덮어쓰므로, 0으로 초기화하는 것은 낭비입니다.

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는 실행 전에 드라이버에 의해 모듈 전역 변수로 설정됩니다. 제어 복제(control replication)는 모든 랭크에서 드라이버를 실행하므로, 모든 워커가 동일한 값을 확인합니다.

각 태스크는 먼저 소비자 뷰를 구축한 후(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. 모든 샤드는 dtype과 후행 축(shape[1:]) 을 공유해야 합니다. 이 레시피는 축 0을 따라 샤드를 쌓습니다. 대상의 후행 축은 trailing_shape에서 비롯되며, 디스커버리 단계에서 이를 첫 번째 파일의 값으로 고정합니다. 샤드별 행 수(shape[0])는 서로 다를 수 있으며, 누적 오프셋 테이블이 이를 처리합니다. 이 예제는 dtype이나 후행 모양이 첫 번째 샤드와 다른 모든 샤드를 설명적인 오류 메시지와 함께 거부합니다.

  2. 변형(variant)과 일치하는 컨슈머를 선택합니다. 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 경로는 단일 노드 데모에서만 작동합니다.

변형

균일 샤드 고속 경로 (파일당 하나의 태스크)

모든 샤드가 이미 동일한 (셰이프, dtype) 을 가지고 있고, 사용 가능한 프로세서 수가 num_shards와 일치하는 경우, cum-rows / bisect 기법은 오버헤드가 됩니다. 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) 가 주어지면, C 연속 영역으로 만들 수 있는 (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 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

(포맷별 내장 단일 호출 로더에 대해서는 이 파일 상단의 “이 스킬이 존재하는 이유” 표를 참조하십시오.)

탐색 단계(1단계)에서는 각 형식의 메타데이터를 파싱합니다. .npy / HDF5 / Parquet은 모두 디스크에 파일별 행 수와 dtype 정보를 포함하고 있습니다. 원시 바이너리 형식은 그렇지 않으므로, 사이드카를 사용하거나 파일 크기에서 파생해야 합니다.

흔히 발생하는 오류

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는 드라이버에서는 문제없지만, 리프 태스크에서는 사용할 수 없습니다. torch 기반의 해결 방법은 examples/dlpack/leaf_task_interop.py를 참조하십시오.

태스크 내 assert는 런타임을 중단시킵니다

Legate는 @task 내에서 발생하지 않은 예외를 계약 위반으로 간주하며, throws_exception()으로 등록된 태스크가 아닌 경우 실행을 중단합니다. 실행 전에 호스트에서 정상성 검사를 수행하십시오.

실행 도메인은 파티션 타일 수와 일치해야 합니다

create_manual_task(launch_shape=...) 와 partition_by_tiling(...)은 서로 독립적이며, 런타임에서는 불일치를 감지하지 않습니다. 실행 도메인이 크면 → 범위 외 타일이 발생하고, 작으면 → 기록되지 않은 타일이 발생합니다. 항상 두 개의 별도 천장( ceil) 연산을 통해 동일한 (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년 6월 29일
jpa-patterns
업데이트 된 시간 2026년 6월 30일
fabric-lakehouse
업데이트 된 시간 2026년 6월 30일
prisma-expert
업데이트 된 시간 2026년 6월 29일
OR