opção
LarLar Skill Gerenciamento de banco de dados cupynumeric-parallel-data-load

cupynumeric-parallel-data-load

NVIDIA/skills NVIDIA/skills

Carregar conjuntos de dados fragmentados (npy, Parquet, HDF5, binário bruto) em matrizes cuPyNumeric distribuídas, utilizando particionamento manual e execução de tarefas pelo Legate.

...Expandir tudo
0
Tempo atualizado 28 de Setembro de 2026

Dados fragmentados em paralelo -> carregamento com cupynumeric

Por que essa habilidade existe. O cupynumeric espelha a API de matrizes do NumPy, incluindo o cupynumeric.load para um único arquivo .npy. Além disso, o carregamento de arquivos ocorre no Legate, não no cupynumeric:

Formato Carregador embutido
Único arquivo .npy cupynumeric.load(path) (compatibilidade com a API do NumPy)
HDF5 (arquivo único) legate.io.hdf5.from_file / from_file_batched
Múltiplos arquivos fragmentados (qualquer formato), Parquet/Arrow, binário bruto, layouts personalizados Não há carregador integrado — esta habilidade.

Esta habilidade mostra a maneira canônica de preencher a lacuna na última linha: escreva uma tarefa Legate em Python que chame o leitor de terceiros exigido pelo formato (h5py, pyarrow, np.memmap, ...) dentro do corpo da tarefa e deixe o Legate distribuir as leituras entre GPUs/nós. Para os formatos com carregador integrado, dê preferência a ele, a menos que você precise de um corpo de tarefa personalizado (carregador baseado em mmap, decodificador específico do formato, metadados sidecar, leituras parciais/fragmentadas).

Padrão canônico: partição manual + lançamento manual da tarefa, dimensionado para a máquina, não para os arquivos. Apenas o eixo 0 é fragmentado; os eixos subsequentes são incluídos dentro de cada bloco. O número de linhas por fragmento pode variar entre os arquivos (apenas o dtype e os eixos subsequentes devem corresponder); o lançamento preenche todos os processadores disponíveis, independentemente do número de arquivos.

.npy é o exemplo prático porque o cabeçalho contém a forma e o tipo de dados no disco, mas a estrutura se aplica a qualquer formato com leituras de intervalo/fatia de baixo custo (binário bruto, HDF5, Parquet/Arrow — veja “Outros formatos” abaixo). Implementação de referência: assets/examples/parallel_npy_load.py.

Suposição sobre o layout dos dados

Esta habilidade trata exclusivamente do carregamento — ela pressupõe que os dados já estejam dispostos em um sistema de arquivos compartilhado de alguma forma previsível e indexável. A geração desses arquivos está fora do escopo (o exemplo inclui um subcomando de gravação por conveniência, mas os usuários reais devem usar os seus próprios).

O exemplo prático pressupõe um layout específico:

  • Um diretório contendo arquivos nomeados shard_0000.npy, shard_0001.npy, ... em uma sequência contígua de números inteiros (largura de 4 preenchida com zeros).
  • Todos os fragmentos compartilham o mesmo dtype e os mesmos eixos finais (shape[1:]); o eixo 0 (linhas por fragmento) pode variar entre os arquivos — a receita constrói uma tabela cumulativa de deslocamento de linha e lê a fatia sobreposta de cada arquivo a partir da tarefa folha.
  • O diretório é visível para todos os ranks (sistema de arquivos compartilhado para execuções com múltiplos nós).

A função discover_layout() do exemplo exibe o que encontrou e apresenta uma falha grave com um erro descritivo quando o layout está incorreto (diretório ausente, ausência de shards, incompatibilidade de dtype/ eixos finais ou uma lacuna na sequência contígua shard_NNNN.npy ).

Se seus dados estiverem em um layout diferente — binário bruto de passo fixo, um arquivo HDF5 com um conjunto de dados por fragmento, uma árvore de diretórios, ... — apenas o padrão glob, o leitor por arquivo (etapa 4 abaixo) e a descoberta de metadados (etapa 1 abaixo) mudam. O mecanismo de particionamento e inicialização é independente do layout.

Quando usar

Consulte a tabela de formatos acima para a decisão de roteamento (carregador embutido vs. esta skill). Além disso, há dois indícios adicionais de que esta skill é a escolha certa:

  • Substituir o código sequencial `np.concatenate([read(f) for f in files]) ` por leituras paralelas por GPU.
  • Demonstrar como uma tarefa Legate Python definida pelo usuário grava em um matriz de saída cupynumeric por meio de uma execução manual.

Exemplos

Os caminhos abaixo são escritos em relação ao diretório desta habilidade (o script vem em assets/examples/parallel_npy_load.py). Ajuste o prefixo para corresponder ao local onde sua habilidade está instalada (por exemplo, skills/cupynumeric-parallel-data-load/assets/... se a habilidade estiver no diretório de nível superior skills/ ).

# Um único nó, 4 GPUs.
legate --gpus 4 --fbmem 4000 --min-gpu-chunk 1 \
    assets/examples/parallel_npy_load.py \
    read --shard-dir /shared/scratch/demo
# Multinó, 2 nós x 4 GPUs (slurm), sistema de arquivos compartilhado em --shard-dir.
# Gere os fragmentos uma vez no rank 0 e, em seguida, execute novamente `read` em qualquer escala.
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

Sem sinalizadores de layout — o driver de leitura percorre cada cabeçalho .npy para recuperar a contagem de linhas por arquivo, a forma final e o tipo de dados (dtype), e então calcula tile_rows a partir da contagem de processadores disponíveis.

--min-gpu-chunk 1 só é necessário quando a contagem de elementos por bloco está abaixo do tamanho mínimo padrão de bloco do Legate para execuções em GPU (por exemplo, os padrões do exemplo prático — total de linhas divididas entre 4 GPUs a cerca de 1 milhão por bloco — ficam abaixo do limite e, caso contrário, seriam agrupados em uma única GPU). Para conjuntos de dados de tamanho de produção (dezenas de milhões de elementos por bloco ou mais), você pode omitir a opção e deixar o Legate usar seu padrão. Aumentá-la para um valor moderado (por exemplo, --min-gpu-chunk 1024) é adequado quando cada bloco é grande o suficiente para que a sobrecarga por tarefa seja mais importante do que atribuir um bloco a cada GPU.

Instruções

Cinco etapas a partir de um exemplo prático em .npy; apenas a etapa 1 (análise do cabeçalho do formato) e a etapa 4 (o leitor por arquivo dentro do corpo da tarefa) são específicas do formato. Os outros três (alocação de destino, partição, barreira) são reutilizados sem alterações em todos os formatos — consulte “Outros formatos” abaixo para os pontos de troca.

1. Leia os metadados de cada fragmento

Faça uma varredura no diretório e examine cada cabeçalho .npy (mmap_mode="r" lê apenas o cabeçalho). O cabeçalho contém a forma e o tipo de dados por fragmento, de modo que o driver pode recuperar o total de linhas, a forma final e uma tabela cumulativa de deslocamento de linhas sem precisar carregar os dados:

paths = sorted(SHARD_DIR.glob("shard_*.npy"))

per_file_rows = []                       # linhas ao longo do eixo 0 por arquivo
trailing_shape = None                    # shape[1:], deve corresponder entre os arquivos
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}: incompatibilidade entre a forma final e o tipo de dados "
            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)  # comprimento N+1
total_rows = int(cum_rows[-1])

O trecho de código acima garante a correspondência entre dtype e trailing_shape (ou seja, shape[1:]) entre os arquivos. A contagem de linhas por fragmento pode variar — a tabela cum_rows lida com isso. O código de produção também deve verificar se os nomes formam uma sequência contígua de shard_0000.npy ... shard_NNNN.npy (omitida do trecho por uma questão de concisão; consulte discover_layout() no exemplo prático). A descoberta depende apenas do que o próprio formato no disco expõe (o cabeçalho .npy aqui, .shape / .dtype para HDF5, etc.); qualquer arquivo auxiliar (manifesto, hashes de conteúdo) é uma etapa de verificação separada adicional.

2. Crie o armazenamento de saída cupynumeric a partir dos metadados

A matriz total abrange total_rows ao longo do eixo 0; os eixos restantes vêm de trailing_shape inalterados. Use cn.empty — a tarefa sobrescreve todas as células, a inicialização com zeros seria um desperdício.

import cupynumeric as cn

total_shape = (total_rows,) + trailing_shape
out = cn.empty(total_shape, dtype=dtype)

3. Divida o armazenamento em blocos de acordo com o número de processadores

A forma de inicialização é dimensionada de acordo com os processadores disponíveis, não com a contagem de arquivos. Escolha tile_rows = ceil(total_rows / num_processors) e particione o eixo 0 por esse tamanho de bloco. Os eixos finais não são particionados (o bloco abrange toda a extensão nesse eixo). É permitido que o último bloco seja curto — é exatamente isso que o `partition_by_tiling` suporta — portanto, a receita não precisa de nenhuma restrição de divisibilidade.

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  # fazer corresponder a contagem de blocos à partição

4. Defina a tarefa folha e inicie-a manualmente

PATHS e CUM_ROWS (os caminhos dos arquivos e a tabela de deslocamento cumulativo de linhas da etapa 1), além de TILE_ROWS, são preenchidos como variáveis globais do módulo pelo driver antes do lançamento; o controle de replicação executa o driver em cada rank, de modo que cada trabalhador veja valores idênticos.

Cada tarefa constrói primeiro sua visão de consumidor (cupy na GPU, numpy na CPU/OMP) e lê a contagem real de linhas do bloco a partir de view.shape[0] — o próprio PhysicalStore não possui o atributo .shape, portanto, é necessário passar pela visão. Em seguida, ela calcula seu intervalo global de linhas a partir de sua coordenada de inicialização e dessa contagem de linhas, divide cum_rows ao meio para os arquivos sobrepostos e copia cada fatia de arquivo sobreposta para a fatia de destino correspondente. Registre as variantes de CPU, OMP e GPU para que a mesma inicialização seja executada sem alterações em qualquer lugar; o despacho por meio de `ctx.get_variant_kind()` seleciona o consumidor que corresponde ao local onde o OutputStore está residente (cp.from_dlpack(dst) para FBMEM, np.asarray(dst) para SYSMEM). O cupy é importado apenas dentro do ramo da GPU, portanto, o corpo da tarefa é carregado em máquinas sem 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]                              # índice do bloco 0..num_tasks-1

    variant = ctx.get_variant_kind()
    if variant == VariantCode.GPU:
        import cupy as cp                              # preguiçoso: somente na GPU
        view = cp.from_dlpack(dst)
    else:
        view = np.asarray(dst)                         # visualização numpy sem cópia

    tile_rows_actual = view.shape[0]                   # falta uma linha no último bloco
    row_start = t * TILE_ROWS                          # início global no eixo 0
    row_end = row_start + tile_rows_actual

    # Encontrar o intervalo semiaberto de índices de arquivo que se sobrepõem a [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):
        # Interseção do bloco [row_start, row_end) com o arquivo [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                # gravação numpy sem cópia

manual_task = runtime.create_manual_task(
    load_tile.library,
    load_tile.task_id,
    (num_tasks,),                                      # número de tarefas lançadas == número de blocos
)
manual_task.add_output(partition)
manual_task.execute()

Ambos os consumidores passam pelos produtores nativos do PhysicalStore (__dlpack__ para cupy, __array_interface__ para np.asarray) — visualizações sem cópia do bloco local. O custo da bissecção é O(log num_shards) e o loop interno normalmente itera de 1 a 2 vezes (os blocos se sobrepõem, no máximo, em alguns arquivos).

5. Fence e verificação

get_legate_runtime().issue_execution_fence(block=True)

Restrições rígidas

  1. Todos os fragmentos devem compartilhar o dtype e os eixos finais (shape[1:]). A receita empilha os fragmentos ao longo do eixo 0; os eixos finais do destino provêm de trailing_shape, que a etapa de descoberta bloqueia no valor do primeiro arquivo. O número de linhas por fragmento (shape[0]) pode variar livremente — a tabela de deslocamento cumulativo lida com isso. O exemplo rejeita qualquer fragmento cujo dtype ou forma final difira do primeiro, exibindo um erro descritivo.

  2. Escolha o consumidor que corresponda à variante. cp.from_dlpack rejeita armazenamentos residentes em SYSMEM; np.asarray retorna silenciosamente uma visão do host de um armazenamento residente em FBMEM no qual você não pode, na verdade, gravar . Faça o despacho com base em ctx.get_variant_kind() para que cada variante utilize seu próprio consumidor — consulte a etapa 4.

  3. As visualizações mmap nem sempre são contíguas em C — envolva cada fatia por arquivo com np.ascontiguousarray(arr[file_lo:file_hi]) antes de .set() ou da gravação in-place do numpy.

  4. Multinó: SHARD_DIR deve estar em um sistema de arquivos compartilhado. Cada trabalhador (em cada rank) abre os fragmentos pelo caminho; caminhos /tmp locais ao nó funcionam apenas para demonstrações de nó único.

Variantes

Caminho rápido com shards uniformes (uma tarefa por arquivo)

Quando cada fragmento já possui a mesma forma e tipo de dados (shape, dtype) e você por acaso tiver num_shards processadores disponíveis, o mecanismo de cum-rows / bisect representa uma sobrecarga. Defina tile_rows = shard_shape[0] e num_tasks = num_shards; a partição passa a ter um bloco por arquivo e cada tarefa lê exatamente um arquivo do início ao fim (sem bisect, sem loop interno). A alteração no driver é feita em uma única linha:

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)

O mesmo corpo da tarefa `load_tile ` continua funcionando em ambos os modos — o loop interno simplesmente itera exatamente uma vez por tarefa. Não há necessidade de um corpo de tarefa separado para o caminho rápido.

Superdecompor para um melhor balanceamento de carga

O padrão tile_rows = ceil(total_rows / num_processors) resulta em um bloco por processador. Para decomposição excessiva por um fator K (blocos menores, mais tarefas pontuais, enfileiramento mais refinado), divida por K * num_processors em vez disso:

tile_rows = max(1, (total_rows + K * num_processors - 1) // (K * num_processors))

num_tasks = ceil(total_rows / tile_rows) então se expande para aproximadamente K * num_processors. O mesmo corpo de tarefa ainda funciona — a bissecção simplesmente resulta em mais tarefas por arquivo.

Outros formatos

Apenas o leitor por arquivo dentro de `load_tile` muda. O contrato do leitor: dado um caminho de arquivo e um intervalo de linhas semiaberto [file_lo, file_hi) ao longo do eixo 0, retorne uma matriz numpy com a forma (file_hi - file_lo,) + trailing_shape que possa ser tornado contíguo em C. São necessárias leituras baratas de intervalos/fatias — formatos que suportam apenas “ler o arquivo inteiro” inviabilizam o caso de sobreposição parcial (um bloco que cobre apenas parte de um arquivo).

Formato Leitor dentro da tarefa folha
.npy (exemplo prático) host = np.ascontiguousarray(np.load(p, mmap_mode="r")[file_lo:file_hi])
Binário bruto (forma fixa) arr = np.memmap(p, dtype=DTYPE, mode="r", shape=(rows_in_file, *trailing_shape)); host = np.ascontiguousarray(arr[file_lo:file_hi])
HDF5 com h5py.File(p, "r") como 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

(Para carregadores embutidos de chamada única por formato, consulte a tabela “Por que esta habilidade existe” no início deste arquivo.)

A etapa de descoberta (etapa 1) analisa os metadados de cada formato: .npy / HDF5 / Parquet contêm, no disco, a contagem de linhas por arquivo + o dtype. O binário bruto não contém — é necessário usar um arquivo sidecar ou derivar a partir do tamanho do arquivo.

Armadilhas comuns

cn.asarray(dst) é inválido em uma tarefa folha

Dentro do corpo de uma @task, qualquer operação cupynumeric que acesse o tempo de execução de nível superior — cn.asarray(store), atribuição de fatia cn_dst[s] = host_np — aciona create_index_space a partir do contexto errado e o Legion é abortado:

EXCEÇÃO DE USO DA API DO LEGION: Contexto de tarefa inválido passado para a chamada de tempo de execução
create_index_space

Solução: utilize a cápsula DLPack com uma biblioteca de terceiros (cupy / torch / numpy) dentro de tarefas folha. cn.asarray funciona normalmente no driver, mas não em tarefas folha. Consulte examples/dlpack/leaf_task_interop.py para ver a solução alternativa com o torch.

Uma assert dentro da tarefa interrompe o tempo de execução

O Legate trata exceções não levantadas em uma @task como uma violação de contrato e interrompe, a menos que a tarefa tenha sido registrada com throws_exception(). Faça uma verificação de integridade no host antes de iniciar.

O domínio de inicialização deve corresponder ao número de blocos de partição

create_manual_task(launch_shape=...) e partition_by_tiling(...) são independentes — o tempo de execução não detecta uma incompatibilidade. Domínio de lançamento maior → blocos fora do intervalo; menor → blocos não gravados. Sempre derive ambos a partir do mesmo (total_rows, tile_rows) por meio de duas divisões por teto separadas (dimensionar o domínio de lançamento diretamente para num_processors resultaria em lançamento excessivo quando 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,))
Ver no 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,))
```

Instalar cupynumeric-parallel-data-load

Baixe e extraia os arquivos de habilidades para o diretório .claude/skills/.

Baixar ZIP

Clone o repositório e copie os arquivos da habilidade para o seu projeto.

git clone https://github.com/NVIDIA/skills/tree/main/skills/cupynumeric-parallel-data-load # Copy SKILL.md to your .claude/skills/ directory

Copiar Copiar
Configuração rápida: Copie a pasta da habilidade para .claude/skills/ O Claude detectará e utilizará a habilidade automaticamente
Repositório NVIDIA/skills

Habilidades relacionadas

microservices-patterns
Tempo atualizado 29 de Junho de 2026
jpa-patterns
Tempo atualizado 30 de Junho de 2026
fabric-lakehouse
Tempo atualizado 30 de Junho de 2026
prisma-expert
Tempo atualizado 29 de Junho de 2026
OR