DALI动态模式Skill dali-dynamic-mode

本文档是NVIDIA DALI命令式动态模式(ndd)的完整指南,涵盖核心数据类型(Tensor/Batch)、有状态读取器用法、设备处理、异步执行模型、线程配置、RNG随机数生成、检查点机制以及从管道模式迁移的详细对照。它帮助AI代理和开发者快速掌握ndd的Python直接调用风格,避免常见错误,高效构建数据加载与GPU处理流程。关键词:DALI动态模式、ndd、NVIDIA DALI、数据加载、数据预处理、GPU加速、命令式API、PyTorch集成、深度学习、检查点、RNG。

GPU数据处理 0 次安装 1 次浏览 更新于 9/7/2026
名称 dali-dynamic-mode
描述 “DALI命令式动态模式(nvidia.dali.experimental.dynamic,即ndd):在编写ndd代码或迁移管道时使用;跳过仅限管道级的任务。”
开源协议 Apache-2.0 metadata:
作者 “DALI Team dali-team@nvidia.com” tags: - dali - 动态模式 - ndd - 数据加载 - 数据处理 - GPU处理 languages: - python team: dali domain: 深度学习

DALI 动态模式

目的

指导AI代理编写、审查和迁移使用DALI命令式动态模式API(nvidia.dali.experimental.dynamic,即ndd)的代码。

说明

  • 将动态模式导入为nvidia.dali.experimental.dynamic as ndd,并以普通Python直接调用ndd的方式编写代码;不要使用管道模式API,如Pipeline@pipeline_defpipe.build()pipe.run()
  • 将有状态读取器视为有状态对象:创建一次,跨多个epoch复用,并将batch_size传递给next_epoch(...)
  • 向随机操作传递显式batch_size;没有管道级batch size可继承。
  • 使用动态模式API约定:device="gpu"而不是管道模式的"mixed",使用Batch.tensors[...]进行样本选择,使用Batch.slice[...]进行逐样本切片。
  • 使用.torch()将张量或批次转换为PyTorch张量。对于形状可变的批次,使用pad=True

先决条件

  • 要运行或验证代码,必须安装支持动态模式导入的NVIDIA DALI,即nvidia.dali.experimental.dynamic
  • GPU解码或GPU运算符需要支持CUDA的DALI构建以及可用的NVIDIA GPU/驱动。
  • 框架转换示例需要安装目标框架,例如PyTorch以支持.torch()

简介

动态模式是DALI的命令式Python API。它允许代码直接从普通的Python控制流调用DALI运算符,而不是构建和运行管道图。

核心数据类型

Tensor – 单个样本

t = ndd.tensor(data)           # 复制
t = ndd.as_tensor(data)        # 包装,尽可能不复制
t.cpu()                        # 移动到CPU
t.gpu()                        # 移动到GPU
t.torch(copy=False)            # 转换为PyTorch张量,默认不复制
t[1:3]                         # 支持切片
np.asarray(t)                  # 通过__array__转换为NumPy(仅限CPU)

支持__dlpack____cuda_array_interface____array__和算术运算符。

Batch – 样本集合(支持可变形状)

b = ndd.batch([arr1, arr2])    # 复制
b = ndd.as_batch(data)         # 包装,尽可能不复制

Batch没有__getitem__batch[i]会引发TypeError,因为索引有歧义(样本选择与逐样本切片)。请使用显式API:

意图 方法 返回
获取样本i batch.tensors[i] Tensor
获取样本子集 batch.tensors[slice_or_list] Batch
在每个样本内切片 batch.slice[...] Batch(保持batch_size不变)
按样本切片 batch.slice[batch_of_indices] Batch(保持batch_size不变)

.tensors[]选择哪些样本.slice在每个样本内部索引。

xy = ndd.random.uniform(batch_size=16, range=[0, 1], shape=2)
crop_x = xy.slice[0]       # 16个标量的Batch,取自每个样本的第一个元素
crop_y = xy.slice[1]       # 16个标量的Batch,取自每个样本的第二个元素
sample_0 = xy.tensors[0]   # Tensor,整个第一个样本[x, y]

高级切片

.slice[] API接受索引的批量,允许用户混合批量和标量值,例如:

imgs = ndd.imread(filenames)  # 如果filenames是列表,则是图像批量
sliced = imgs.slice[
    42 :  # 范围起始广播到所有样本
    ndd.batch(imgs.shape).slice[0] // 2  # 逐样本范围截止(每张图像的一半)
]

PyTorch转换:

  • batch.torch() – 适用于均匀形状;对形状参差不齐的批次会引发错误
  • batch.torch(pad=True) – 将参差不齐的批次零填充到最大形状(用于可变长度音频、检测框等)
  • batch.torch(copy=None)是默认值(尽可能避免复制)
  • Batch没有__dlpack__ – 对于DLPack消费者,请先使用ndd.as_tensor(batch)ndd.as_tensor也支持pad
  • Tensor.torch(copy=False)是默认值(不复制)

迭代:for sample in batch:产生Tensor。

读取器

读取器是有状态对象 – 创建一次,跨epoch复用。这很重要,因为读取器维护内部状态,如洗牌顺序和分片位置。

reader = ndd.readers.File(file_root=image_dir, random_shuffle=True)

for epoch in range(num_epochs):
    for jpegs, labels in reader.next_epoch(batch_size=64):
        # jpegs, labels是Batch对象
        ...

要点:

  • 读取器输出(jpegs, labels等)是CPU张量/批次。标签通常保留在CPU上,直到转换为框架格式(例如labels.torch().to(device))。
  • 读取器类是PascalCasendd.readers.File(...)ndd.readers.COCO(...)ndd.readers.TFRecord(...)
  • batch_size传递给next_epoch(),而不是读取器构造函数
  • 不带batch_size的next_epoch(batch_size=N)产生Batch元组;不带batch_size的next_epoch()产生Tensor元组
  • 迭代器必须完整消费后才能再次调用next_epoch()
  • 一旦读取器以特定batch_size使用,就不能更改。同样,以批次模式使用的读取器不能切换到样本模式,反之亦然。

分布式训练的分片读取:

reader = ndd.readers.File(
    file_root=image_dir,
    shard_id=rank, num_shards=world_size,
    stick_to_shard=True,
    pad_last_batch=True,
)

设备处理

  • 设备从输入推断 – 如果任何输入在GPU上,则使用GPU
  • 对于混合解码:使用device="gpu"(不是"mixed")。"mixed"关键字是管道模式概念,用于隐式CPU到GPU传输;在动态模式下,传递device="gpu"会触发相同的硬件加速解码路径。
  • 在传递给GPU模型之前不要调用.cpu().torch()直接给出GPU张量。.cpu()仅对需要主机内存的消费者(numpy、__array__)是必需的。
  • DALI和PyTorch之间的CUDA流同步通过DLPack自动完成 – 无需手动流管理。

执行模型

默认模式为eager – 在后台线程异步执行,立即返回。

**在大多数情况下无需.evaluate()。**任何数据消费(.torch()__dlpack____array__.shape、属性访问、迭代)都会自动触发求值。

对于调试,切换到同步模式,使错误在精确调用点显现,而不是稍后在异步队列中出现:

with ndd.EvalMode.sync_cpu:
    images = ndd.decoders.image(jpegs, device="gpu")
    images = ndd.resize(images, size=[224, 224])
    # 任何错误都在此处显现,在失败的操作处

模式(同步性增加):deferred < eager < sync_cpu < sync_full

使用EvalMode.sync_full进行调试,而不是散布.evaluate()调用-- 更简洁,并一次性捕获所有问题。sync_cpu通常足够且比sync_full更轻量。

线程配置

ndd.set_num_threads(4)  # 在启动时调用一次,仅在需要覆盖默认值时使用

控制DALI用于CPU运算符的内部工作线程数。默认为CPU亲和计数或DALI_NUM_THREADS环境变量。与Python级线程无关。

RNG

两种方法(使用一种,不要同时使用):

# 方法1:设置线程本地默认种子(简单,适用于大多数情况)
ndd.random.set_seed(42)
angles = ndd.random.uniform(batch_size=64, range=(-30, 30))

# 方法2:显式RNG对象(更精细控制,将rng=传递给每个操作)
rng = ndd.random.RNG(seed=42)
values = ndd.random.uniform(batch_size=64, range=[0, 1], shape=2, rng=rng)

rng=传递给随机操作时,显式RNG覆盖默认种子。线程本地:每个线程有独立的随机状态。

随机操作在处理批次时需要显式batch_size – 没有管道级batch size可继承。

检查点

动态模式没有管道级检查点。检查点聚合各个有状态对象的状态:读取器和RNG实例。无状态操作(解码器、resize、rotate、normalize等)不属于检查点的一部分。

ckpt = ndd.checkpoint.Checkpoint()
ckpt.register(reader, "my_reader")
ckpt.register(rng, "rng")

# ... 迭代一段时间 ...

ckpt.collect()                       # 快照已注册对象
ckpt.save("ckpt_{seq:04d}.json")     # 写入ckpt_0000.json, ckpt_0001.json, ...

恢复是对称操作 – 构建一个新的读取器和RNG,然后load + register。加载的状态在register时应用于每个对象:

reader = ndd.readers.File(file_root=..., enable_checkpointing=True, name="my_reader")
rng = ndd.random.RNG()

ckpt = ndd.checkpoint.Checkpoint()
ckpt.load("ckpt_{seq:04d}.json")     # 选择最高序号
ckpt.register(reader, "my_reader")   # 在此处应用状态
ckpt.register(rng, "rng")            # 同上

for batch in reader.next_epoch(batch_size=N):
    ...  # 产生检查点迭代之后的下一个批次

关键规则:

  • 读取器必须选择加入。 使用enable_checkpointing=True构造。注册已经迭代过的读取器而没有此标志会引发RuntimeError;如果读取器尚未迭代,则register会追溯启用它。
  • 读取器状态必须在第一次next_epoch调用之前应用。 预取线程在第一次迭代时启动,快照队列在此之后锁定。对已迭代的读取器调用set_state(或从加载的检查点注册)会引发RuntimeError
  • enable_checkpointing=Truecompile=True不兼容。 在启用检查点的读取器上调用reader.next_epoch(..., compile=True)会引发NotImplementedError
  • 命名注册更安全。 匿名register(op)使用顺序键(__op_0__op_1,…),因此保存和恢复之间的注册顺序必须匹配。类型标记可以捕获跨类型交换,但不能捕获兼容类型的重排序。首选register(op, name)
  • ndd.checkpoint.current() 返回绑定到当前线程本地EvalContextCheckpoint。跨调用共享 – 如果重用默认上下文进行不相关的运行,请调用ckpt.clear()
  • 文件名模式: save/load接受带单个{seq}占位符的Python格式字符串(例如"ckpt_{seq:04d}.json")。save选择下一个空闲序号;load选择磁盘上匹配的最高序号。
  • 格式版本严格。 如果载荷来自不同的检查点格式版本,deserialize会拒绝 – 无自动升级。
  • 非线程安全。 每个线程一个Checkpoint

每个ReaderRNG上也可以直接使用手动get_state/set_stateCheckpoint聚合器基于此构建。仅在与外部检查点系统集成时使用手动API。

示例

图像分类管道

import nvidia.dali.experimental.dynamic as ndd

reader = ndd.readers.File(file_root="/data/imagenet/train", random_shuffle=True)

for epoch in range(num_epochs):
    for jpegs, labels in reader.next_epoch(batch_size=64):
        images = ndd.decoders.image(jpegs, device="gpu")
        images = ndd.resize(images, size=[224, 224])
        images = ndd.crop_mirror_normalize(
            images,
            mean=[0.485 * 255, 0.456 * 255, 0.406 * 255],
            std=[0.229 * 255, 0.224 * 255, 0.225 * 255],
        )
        train_step(images.torch(), labels.torch())

常见错误

错误 正确 原因
device="mixed" device="gpu" "mixed"仅用于管道模式
batch[i] batch.tensors[i] Batch没有__getitem__
batch.tensors[0]用于逐样本切片 batch.slice[0] .tensors选择样本;.slice在每个样本内切片
每个操作后调用.evaluate() 让消费触发求值 .torch().shape等会自动触发
GPU模型前调用.cpu() 直接.torch() 避免浪费D2H + H2D往返
每个epoch重新创建读取器 reader.next_epoch() 读取器有状态 – 创建一次,重用
ndd.readers.file(...) ndd.readers.File(...) 读取器类是PascalCase
next_epoch()循环break 耗尽迭代器或创建新读取器 迭代器必须完整消费后才能再次调用next_epoch()
随机操作未给batch_size ndd.random.uniform(batch_size=N, ...) 没有管道级batch size可继承
在第一次next_epochregister(reader)以恢复 在第一次迭代前注册新构建的读取器 读取器状态只能在预取线程启动前应用
在迭代后恢复到未用enable_checkpointing=True构建的读取器 在构造时传递enable_checkpointing=True(或在第一次迭代前注册) 否则后端不保留快照
显式写出默认参数值 跳过默认参数值 Python侧开销非常高,特别是当参数接受Tensor/Batch时。跳过参数使用快速路径,实际传递哨兵值。

管道模式迁移

管道模式 动态模式
@pipeline_def / pipe.build() / pipe.run() 循环中的直接函数调用
fn.readers.file(...) ndd.readers.File(...)(PascalCase,有状态)
fn.decoders.image(jpegs, device="mixed") ndd.decoders.image(jpegs, device="gpu")
fn.op_name(...) ndd.op_name(...)
管道级batch_size=64 reader.next_epoch(batch_size=64) + 随机操作batch_size=64
管道级seed=42 ndd.random.set_seed(42)ndd.random.RNG(seed=42)
管道级num_threads=4 启动时ndd.set_num_threads(4)
output.at(i) batch.tensors[i]
output.as_cpu() batch.cpu()
pipe.run()返回TensorList元组 reader.next_epoch(batch_size=N)产生Batch元组
Pipeline(..., enable_checkpointing=True) + pipe.checkpoint() / pipeline(checkpoint=...) ndd.checkpoint.Checkpoint + 每个对象register / collect / save / load;读取器选择加入enable_checkpointing=True

局限

动态模式比管道模式更灵活,但性能可能稍差。为获得最大吞吐量,请优先使用管道模式。

故障排除

  • 如果错误晚于失败调用才出现,请在EvalMode.sync_cpuEvalMode.sync_full下重新运行块。
  • 如果读取器跨epoch行为异常,请检查它是否只创建一次,并确保每个next_epoch()迭代器被完整消费。