agent_loop.py — 实现了一步离策略专用的 Agent 循环管理器¶
文件路径: verl/experimental/one_step_off_policy/agent_loop/agent_loop.py
文件概述¶
实现了一步离策略专用的 Agent 循环管理器。与基类 AgentLoopManager 的主要区别是提供了异步版本的 generate_sequences_async,使得推理可以与训练并行执行。
关键代码讲解¶
OneStepOffAgentLoopManager¶
class OneStepOffAgentLoopManager(AgentLoopManager):
async def generate_sequences_async(self, prompts: DataProto) -> DataProto:
"""异步版本的序列生成"""
# 将 batch 分割分发给各 Worker
chunks = prompts.chunk(len(self.agent_loop_workers))
# 使用 asyncio.gather 并行调度所有 Worker
outputs = await asyncio.gather(
*[
asyncio.to_thread(ray.get, worker.generate_sequences.remote(chunk))
for worker, chunk in zip(self.agent_loop_workers, chunks, strict=True)
]
)
# 合并结果
output = DataProto.concat(outputs)
# 计算性能指标
metrics = [output.meta_info.pop("metrics") for output in outputs]
timing = self._performance_metrics(metrics, output)
output.meta_info = {"timing": timing, **outputs[0].meta_info}
return output
关键设计:asyncio.to_thread(ray.get, ...) 将阻塞的 ray.get 调用放到线程池中执行,使得主事件循环不被阻塞。这样训练步骤可以与推理步骤在同一个 asyncio 事件循环中交替执行。
辅助方法¶
async def wake_up(self):
"""唤醒所有推理副本"""
await asyncio.gather(*[replica.wake_up() for replica in self.rollout_replicas])
async def sleep(self):
"""休眠所有推理副本"""
await asyncio.gather(*[replica.sleep() for replica in self.rollout_replicas])
async def clear_kv_cache(self):
"""清除所有推理副本的 KV 缓存"""
await asyncio.gather(*[replica.clear_kv_cache() for replica in self.rollout_replicas])
核心类/函数列表¶
| 名称 | 类型 | 说明 |
|---|---|---|
OneStepOffAgentLoopManager |
类 | 一步离策略 Agent 循环管理器 |
generate_sequences_async() |
异步方法 | 异步序列生成 |
wake_up() / sleep() |
异步方法 | 推理副本电源管理 |
clear_kv_cache() |
异步方法 | 清除 KV 缓存 |
与其他模块的关系¶
- 继承自
AgentLoopManager(agent_loop/agent_loop.py) - 被
OneStepOffRayTrainer(ray_trainer.py)使用 - 管理标准的
AgentLoopWorker(不需要部分回滚支持)
小结¶
OneStepOffAgentLoopManager 的核心贡献是 generate_sequences_async 方法。它通过 asyncio.to_thread 将 Ray 的同步调用转为异步,实现了与训练步骤的流水线重叠。相比 FullyAsyncAgentLoopManager,它更简单——不需要取消/恢复机制,因为一步离策略模式下推理总是在训练之前完成。