
TransferQueue 的关键不是把 DataProto 换成另一个容器,而是把“函数返回完整 batch”改成“产物写入外部队列,trainer 用 metadata 取数”。
第 27 篇把 data movement 拆成 DataProto序列化、driver 侧合并、Ray object 和 DataProtoFuture。第 28 篇继续往前走:如果 controller 上的数据回流已经成为结构性压力,系统能不能不要让每个 rollout 结果都先回到 controller,再由 controller 拼出完整 batch?
本文的核心判断是:TransferQueue 路线把同步 PPO 的数据合同从“controller 持有完整 DataProto”改成“TransferQueue 持有真实样本字段,ReplayBuffer 持有可采样 metadata,trainer 只携带 KVBatchMeta”。这让 agent loop 可以把每个完成的轨迹写入队列,让 trainer 按 key 和字段读取后续阶段需要的数据;但代价是系统必须管理 status、key、tag、padding、清理和队列一致性。它是从单 controller dataflow 走向流式数据系统的一步,还不是第 29 篇要讲的 fully async policy。
先看整体路线。读图时注意,数据不再沿着一条 rollout -> DataProto -> trainer的返回路径走,而是分成两条线:真实字段进入 TransferQueue,trainer 拿到的是可定位这些字段的 KVBatchMeta。

TransferQueue 把同步返回改成队列化数据流
这张图解释了第 28 篇和第 27 篇的关系:第 27 篇讨论 DataProto 过 controller 的成本,第 28 篇讨论如何把这部分回流改造成外部存储和 metadata 驱动。源码上,这条路线不是默认开启;ppo_trainer.yaml里 transfer_queue.enable默认是 False,main_ppo_sync.main()会显式把它设为 True,main_ppo.py在开启时给 Ray runtime env 写入 TRANSFER_QUEUE_ENABLE=1(verl/trainer/config/ppo_trainer.yaml:310-320,verl/trainer/main_ppo_sync.py:1797-1806,verl/trainer/main_ppo.py:67-76)。
普通同步路径里,trainer 调用 async_rollout_manager.generate_sequences()后等待一个完整 DataProto 返回。TransferQueue 路线改掉的是这个返回模型:manager 把 prompt 分给 agent loop workers,worker 立刻为每条样本创建后台任务,真正的 agent loop 输出完成后写入 TransferQueue。
下面这张图画的是 producer 侧。读图时注意两层 status:manager 先把 prompt 标记成 running,worker 完成后把每条输出按 {uid}_{session_id}_{index}写入队列,并用 tag 标记 success、长度和 step 信息。

Agent loop 输出如何写入 TransferQueue
源码上,AgentLoopManagerTQ.generate_sequences()先把当前 prompt 的 uid注册到 replay buffer,status 是 running;随后把 TensorDict chunk 给多个 AgentLoopWorkerTQ,这里只等待远端方法被调度,不等待每条 agent loop 轨迹完成(verl/trainer/main_ppo_sync.py:463-483)。worker 侧 generate_sequences()会为 batch 中每条样本创建 asyncio.create_task(),后台执行 _run_prompt()(verl/trainer/main_ppo_sync.py:292-340)。
真正写队列发生在 _agent_loop_postprocess()。它把一个或多个 AgentLoopOutput转成 TensorDict 字段,包括 input_ids、position_ids、multi_modal_inputs、loss_mask、reward 相关字段和轨迹字段;然后调用 tq.async_kv_batch_put(),key 格式是 {uid}_{session_id}_{index},tag 里带 global_steps、status=success、prompt_len、response_len和 seq_len(verl/trainer/main_ppo_sync.py:367-437)。这就是“流式”的第一层含义:产物可以按轨迹写入外部队列,而不是等所有输出拼成一个返回值。
队列里有真实字段,但 trainer 不应该频繁扫描完整数据。ReplayBuffer的角色是维护一个轻量 metadata 视图:它周期性调用 tq.kv_list(),把不同 partition 下的 key 和 tag 放进内存索引。
下面这张图展示 consumer 侧的采样。读图时注意:ReplayBuffer 返回的是 KVBatchMeta,它包含 partition、keys、tags 等定位信息,不是完整 TensorDict。

ReplayBuffer 用 metadata 采样 KVBatchMeta
ReplayBuffer.__init__()会启动后台 poll 线程,_poll_from_transfer_queue()周期性执行 tq.kv_list(),把结果 add 到 self.partitions(verl/trainer/main_ppo_sync.py:193-220)。sample()按 partition_id和 global_steps查 key:如果遇到 running,继续等待;如果看到 success,把 key 和 tag 收集起来;最后返回 KVBatchMeta(partition_id, keys, tags)(verl/trainer/main_ppo_sync.py:249-282)。
这个设计有一个重要边界:第 28 篇的 TransferQueue 路线仍然在 step()里等待当前 global_steps的样本完成。PPOTrainer.step()先调用 generate_sequences(batch),再在 marked_timer("gen")里调用 replay_buffer.sample(partition_id="train", global_steps=self.global_steps)(verl/trainer/main_ppo_sync.py:1637-1661)。所以它解决的是数据回流和存储形态,不是完全取消 trainer/rollout 的同步关系。样本新鲜度和 fully async 的问题留到第 29 篇。
拿到 KVBatchMeta后,后续阶段不需要立刻把完整样本拉回 driver。它们通常按当前阶段需要的字段读取,计算新字段,再写回 TransferQueue。
下面这张图把 old logprob、ref/value、adv 和 metrics 放在同一条线上。读图时注意:每段都围绕同一组 key 工作,但字段集合不同,这就是按字段数据系统和完整 DataProto 回流的差别。

后续阶段按字段读写 TransferQueue
_compute_old_log_prob()先把KVBatchMeta交给 actor worker 计算 logprob;随后只取entropy、log_probs、response_mask等字段,把 nested logprob 转成 padding 形态,再通过tq.kv_batch_put()写回old_log_probs和entropy(verl/trainer/main_ppo_sync.py:1256-1320)。reference 和 critic values 也是类似模式:worker 计算后,从队列取log_probs或values与response_mask,转换后写回ref_log_prob或values(verl/trainer/main_ppo_sync.py:1323-1375)。
advantage 阶段必须在 driver 上做更多算法计算,所以它会取 uid、response_mask、reward、old/ref logprob、values 等字段,临时构造 DataProto 计算 advantage,再把 advantages、returns和可选 rollout correction 字段写回队列(verl/trainer/main_ppo_sync.py:1377-1442)。最后 _compute_metrics()才把 metrics 需要的字段取回并转成 DataProto,用于 compute_data_metrics()、compute_timing_metrics()和 throughput 计算(verl/trainer/main_ppo_sync.py:1498-1533)。
这带来一个清晰的工程解释:TransferQueue 不消灭所有 data movement,但把“整批数据每阶段回到 controller”变成“按 key 和字段选择性读写”。这特别适合多轮 agent、可变输出数、对象列较重、response 长短差异大的场景。
tqbridge把 worker 方法接到队列语义如果只有 trainer 侧手写 kv_batch_get()和 kv_batch_put(),TransferQueue 会变成一堆特殊代码。verl 用 tqbridge()把这件事接进原有 @register和 WorkerGroup 体系。
下面这张图要看的重点是 worker 边界:worker 方法收到 KVBatchMeta时,桥接层先把 metadata 转成真实 TensorDict;方法返回 TensorDict 时,桥接层再把输出写回队列,并返回新的 metadata。

tqbridge 如何连接 KVBatchMeta 和 worker TensorDict
tqbridge()会查参数里是否有 BatchMeta或 KVBatchMeta。如果有,它先初始化 TransferQueue,把 KVBatchMeta转成 BatchMeta,再通过 _meta_to_realdata()或 _async_meta_to_realdata()从队列取真实 TensorDict;worker 函数执行后,如果输出是有 batch size 的 TensorDict,就通过 _update_meta_with_output()把输出字段写回队列(verl/utils/transferqueue_utils.py:298-370、372-424)。
BatchData也加入了 KVBatchMeta和 BatchMeta支持:chunk()遇到 KVBatchMeta会先转换成 BatchMeta,避免每个 rank 频繁在 PUT/GET 时和 controller 通信;concat()遇到 BatchMeta会 concat 后再转回 KVBatchMeta(verl/protocol.py:1231-1324)。这说明 TransferQueue 不是绕开 WorkerGroup,而是把 metadata 形态嵌进原来的 dispatch/collect 体系。
TransferQueue 路线买到的是三件事。第一,rollout producer 不必把完整 DataProto 返回给 controller;第二,trainer 可以用 KVBatchMeta做阶段编排,只在需要时按字段取数;第三,多输出 agent loop 和不同 prompt 的动态 n更自然,因为 key 可以表达 {uid}_{session_id}_{index}。
但它也引入新的系统账。ReplayBuffer.sample()要维护 running/success/failure状态语义;训练结束后要用 tq.kv_clear()清理队列,并从 replay buffer remove key(verl/trainer/main_ppo_sync.py:1625-1627)。当可变轨迹数导致 batch 不整除 DP 或 mini-batch 时,upsample_batch_to_divisible_size()会往 TransferQueue 里写入带 is_padding=True的 synthetic sample,后续 metrics 还要过滤 padding(verl/trainer/ppo/padding_utils.py:127-198,verl/trainer/main_ppo_sync.py:1498-1524)。
所以第 28 篇的结论要保守:TransferQueue 是从单 controller dataflow 向流式数据系统过渡的路线,它把 data movement 从“返回完整 batch”改成“外部队列 + metadata + 字段读写”。这可以缓解 controller 数据回流压力,也让多轨迹 agent 训练更自然;但它需要额外的队列生命周期、状态管理、padding、清理和一致性约束。
放回系列地图,第 27 篇告诉我们 data movement 会藏在阶段之间;第 28 篇则给出一条工程路线:把 rollout 产物写入 TransferQueue,让 ReplayBuffer 用 metadata 组织样本,让后续 worker 和 trainer 按 key/field 读写。
这一步已经不像单机训练循环,也不只是 Ray RPC。它开始接近生产数据系统:有 producer、consumer、metadata index、storage backend、status、cleanup 和 backpressure 风险。下一篇第 29 篇会继续推进到 fully async policy:当 trainer 不再只等待当前 step 的样本,而是从持续生产的样本队列里凑 batch 时,吞吐、新鲜度和 off-policy correction 会形成新的三角关系。
verl/trainer/config/ppo_trainer.yaml:310-320、verl/trainer/main_ppo.py:67-76、verl/trainer/main_ppo_sync.py:1797-1806:TransferQueue 配置入口、Ray runtime env 和同步 TQ runner 启动路径。verl/trainer/main_ppo_sync.py:193-282:ReplayBuffer如何 poll tq.kv_list(),并按 global_steps返回 KVBatchMeta。verl/trainer/main_ppo_sync.py:292-437:AgentLoopWorkerTQ如何创建后台任务,并把 agent loop 输出写入 TransferQueue。verl/trainer/main_ppo_sync.py:463-483:AgentLoopManagerTQ.generate_sequences()如何标记 running、chunk prompt 并触发 worker 生成。verl/trainer/main_ppo_sync.py:1637-1702:TQ 版 PPO step 如何从启动生成、ReplayBuffer 采样走到 old/ref/value/adv/update。verl/trainer/main_ppo_sync.py:1256-1442:old logprob、ref logprob、values 和 advantage 如何按字段读写 TransferQueue。verl/trainer/main_ppo_sync.py:1498-1533、1625-1627:metrics 阶段如何取回必要字段,以及 step 结束后的队列清理。verl/utils/transferqueue_utils.py:253-295:KVBatchMeta与 BatchMeta如何互转。verl/utils/transferqueue_utils.py:298-424:tqbridge()如何把 metadata 输入 materialize 成 TensorDict,并把 worker 输出写回 TransferQueue。verl/protocol.py:1231-1324:BatchData对 KVBatchMeta、BatchMeta的 chunk/concat 支持。verl/trainer/ppo/padding_utils.py:127-198:可变轨迹数下如何向 TransferQueue 写入 synthetic padding sample。