| 名称 | 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_def、pipe.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))。 - 读取器类是PascalCase:
ndd.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=True与compile=True不兼容。 在启用检查点的读取器上调用reader.next_epoch(..., compile=True)会引发NotImplementedError。- 命名注册更安全。 匿名
register(op)使用顺序键(__op_0,__op_1,…),因此保存和恢复之间的注册顺序必须匹配。类型标记可以捕获跨类型交换,但不能捕获兼容类型的重排序。首选register(op, name)。 ndd.checkpoint.current()返回绑定到当前线程本地EvalContext的Checkpoint。跨调用共享 – 如果重用默认上下文进行不相关的运行,请调用ckpt.clear()。- 文件名模式:
save/load接受带单个{seq}占位符的Python格式字符串(例如"ckpt_{seq:04d}.json")。save选择下一个空闲序号;load选择磁盘上匹配的最高序号。 - 格式版本严格。 如果载荷来自不同的检查点格式版本,
deserialize会拒绝 – 无自动升级。 - 非线程安全。 每个线程一个
Checkpoint。
每个Reader和RNG上也可以直接使用手动get_state/set_state – Checkpoint聚合器基于此构建。仅在与外部检查点系统集成时使用手动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_epoch后register(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_cpu或EvalMode.sync_full下重新运行块。 - 如果读取器跨epoch行为异常,请检查它是否只创建一次,并确保每个
next_epoch()迭代器被完整消费。