乐于分享
好东西不私藏

VERL 源码精读 05:DataProto、TensorDict 与跨进程数据协议

VERL 源码精读 05:DataProto、TensorDict 与跨进程数据协议

学习主题:VERL 源码精读 05 

深入理解 DataProto 为什么存在,以及它如何支撑 split、chunk、union、concat、to(device)、Ray 序列化。

 阅读目标:建立可复述、可定位、可调试的源码理解,不停留在 API 名称层面。

1. 本篇学习范围

深入理解 DataProto 为什么存在,以及它如何支撑 split、chunk、union、concat、to(device)、Ray 序列化。

这一篇关注的不是“怎么把命令跑起来”,而是读清楚源码中每个对象为什么存在、它和前后模块如何传递数据,以及出问题时应该从哪里开始定位。

2. 源码入口文件

本篇建议按下面顺序阅读:

verl/protocol.pytests/test_protocol_on_cpu.pytests/test_protocol_v2_on_cpu.pyverl/utils/transferqueue_utils.pydocs/data/transfer_queue.md

这些文件覆盖了本篇主题的主路径。阅读时不需要一开始就把所有分支展开,先抓主干,再回头看特殊配置、后端差异和异常处理。

3. 核心调用链

本篇主调用链可以压缩为:

`DataProto` 由 `batch: TensorDict`、`non_tensor_batch: dict`、`meta_info: dict` 三部分组成。tensor 字段进入 TensorDict,字符串、对象、reward 元信息等进入 non_tensor_batch。controller 通过 `chunk()` 或 `split()` 将 DataProto 拆给多个 worker。worker 计算完成后通过 `DataProto.concat()` 或 `union()` 合并结果字段。Ray 传输时 `__getstate__` 和 `__setstate__` 处理 TensorDict 序列化与恢复。

这条链路是读源码时的地图。后续遇到新的类、函数或配置字段,都可以先判断它属于链路中的哪一段。

4. 核心概念

  • DataProto 是 VERL 各模块之间的最小共同语言。
  • batch
     要求 batch 维一致,适合存放 input_ids、responses、log_probs、advantages。
  • non_tensor_batch
     同样具有 batch 维,只是元素可以是字符串、dict 或对象。
  • meta_info
     不一定具有 batch 维,通常保存全局采样参数、step、统计信息。

这些概念共同决定了本篇源码的设计方式。VERL 的一个重要特点是,很多类名看起来像普通工程封装,但背后实际是在解决大模型 RL 的分布式执行问题。

5. 数据在这一层如何流动

在 VERL 中,几乎所有训练阶段都可以用同一套数据流语言描述:

输入 batch-> 按并行度或任务类型拆分-> 分发到本地函数或远程 worker-> 执行高成本计算或轻量控制逻辑-> 收集结果-> 写回 DataProto / TensorDict / TransferQueue-> 进入下一阶段
Controller(Actor)    → dp.chunk(N) 切分DataProto 【protocol.py】    → ray.remote 调用worker方法    → Ray pickle对象 → 触发 DataProto.__getstate__() 【protocol.py】    → 对象进入Ray ObjectStore,网络传输Worker(Actor)收到对象    → Ray unpickle → 触发 DataProto.__setstate__() 重建DataProto 【protocol.py】    → worker做计算,返回DataProto分片    → 再次触发 __getstate__ 传回controllerController侧ray.get拿到结果    → __setstate__重建各个worker的DataProto分片    → DataProto.concat([dp1,dp2,...]) 合并分片

本篇主题对应的数据流重点是:

  1. DataProto
     由 batch: TensorDictnon_tensor_batch: dictmeta_info: dict 三部分组成。
  2. tensor 字段进入 TensorDict,字符串、对象、reward 元信息等进入 non_tensor_batch。
  3. controller 通过 chunk() 或 split() 将 DataProto 拆给多个 worker。
  4. worker 计算完成后通过 DataProto.concat() 或 union() 合并结果字段。
  5. Ray 传输时 __getstate__ 和 __setstate__ 处理 TensorDict 序列化与恢复。

读代码时要始终跟踪字段而不是只跟踪函数名。典型字段包括:

promptsresponsesattention_maskposition_idsold_log_probsref_log_probvaluesrm_scorestoken_level_rewardsadvantagesreturnsmetrics

不是每一天都会出现全部字段,但这些字段构成了 VERL PPO/GRPO 训练的共同词汇表。

6. 和前后模块的关系

本篇主题通常不是孤立工作的。它至少会连接三类模块:

上游:  配置、数据、controller 状态、已有 batch 字段。本层:  当前主题负责的调度、计算、转换或封装逻辑。下游:  rollout、reward、advantage、loss、metrics、checkpoint 或异步队列。

因此读源码时要避免只看单个函数。更稳的方式是:

1. 找到谁调用它。2. 找到它读取哪些字段。3. 找到它写出哪些字段。4. 找到这些字段下一步被谁使用。5. 找到配置项如何改变它的分支。

这个五步法适用于 VERL 的大多数文件。

7. 实现细节拆解

本篇源码中最值得关注的细节包括:

  1. DataProto
     由 batch: TensorDictnon_tensor_batch: dictmeta_info: dict 三部分组成。
  2. tensor 字段进入 TensorDict,字符串、对象、reward 元信息等进入 non_tensor_batch。
    以from_single_dict作为示例,因为这个最清晰能看见两种字段的区别。进入from_dict后:
  3. controller 通过 chunk() 或 split() 将 DataProto 拆给多个 worker。
  4. worker 计算完成后通过 DataProto.concat() 或 union() 合并结果字段。
    controller: batch(DataProto)    → dispatch_dp_compute_data_proto: batch.chunk(N) 切分成N份 → 分发N个workerworker_i: 拿到分片,计算,返回分片DataProto_icontroller收集全部worker输出 → collect_dp_compute_data_proto → DataProto.concat([dp_1, dp_2,...dp_N]) → 恢复完整大batch
  5. Ray 传输时 __getstate__ 和 __setstate__ 处理 TensorDict 序列化与恢复。
    序列化(pickle / Ray 发送对象时调用)
两种序列化方式,一是numpy模式:调用 serialize_tensordict → TensorDict转numpy dict;或者默认模式:用 torch.save 把TensorDict直接dump进BytesIO,返回bytes。
返回的 tuple (batch_payload, non_tensor_batch, meta_info) 才是真正被 pickle/Ray 传输的内容;不再直接 pickle 原始 TensorDict 对象
__setstate__则是反序列化(Ray 接收对象 /pickle.load 时调用)

这些步骤背后通常有两类逻辑:

控制逻辑:  判断当前训练需要哪些角色、哪些字段、哪些分支。计算逻辑:  真正执行模型推理、训练、reward、advantage 或 loss。

HybridFlow 的设计要求我们把这两类逻辑区分开。控制逻辑更适合在 trainer 或 manager 中读;计算逻辑更适合在 worker、engine 或 core_algos 中读。

8. 配置如何影响本篇路径

VERL 的同一段源码经常会被配置切到不同路径。阅读本篇时尤其要关注这些配置类型:

algorithm:  决定 PPO、GRPO、KL、advantage、rollout correction 等算法行为。actor_rollout_ref:  决定 actor、rollout、reference policy、model path、训练后端和推理后端。critic:  决定是否启用 value model,以及 critic 的训练后端。reward:  决定使用规则 reward、reward model、remote reward 还是 sandbox reward。trainer:  决定训练步数、资源规模、logger、validation、checkpoint 和 V1/V0 模式。

读源码前最好先打印 resolved config。否则很容易在一个未启用的分支里浪费时间。

9. 常见误区和调试要点

  • 新增字段时要判断它是 tensor、非 tensor,还是 meta 信息;放错位置会影响切分和合并。
  • 两个 DataProto 做 union 时,同名字段必须一致,否则会触发一致性检查。
  • Ray 传输大 batch 时,DataProto 序列化开销可能成为性能瓶颈。

调试 VERL 时,不建议一上来就改源码。更稳的顺序是:

1. 确认命令行 override 是否真的进入 resolved config。2. 确认当前走 V0 还是 V1 trainer。3. 确认 DataProto / TensorDict 里字段是否存在。4. 确认 batch 维、response 长度、mask 是否一致。5. 确认对应 worker 是否真的被创建。6. 确认 Ray worker 日志里的原始异常。7. 最后再判断是不是算法公式或 loss 本身的问题。

这个顺序能避免把配置错误、数据错误、分布式调度错误误判成算法错误。

10. 本篇压缩总结

读懂这一篇后,应该能够回答三个问题:

1. 这一层在 VERL 训练链路中负责什么?2. 它接收哪些字段,又产出哪些字段?3. 它的行为主要由哪些配置项改变?

如果这三个问题能答清楚,就说明不是在背目录,而是在按 VERL 的真实执行路径读源码。

11. 下一步阅读

读完本篇后,建议继续沿着训练链路向后走:

入口与配置-> 数据和 DataProto-> WorkerGroup 和 Ray 调度-> Worker 与模型引擎-> Rollout-> Reward-> Advantage-> Actor/Critic loss-> Metrics、Checkpoint、异步扩展

VERL 源码量很大,但主线并不乱。只要始终围绕“一个 batch 如何从 prompt 变成 actor update”这条线阅读,就能把分散目录组织成一张完整图。