执行摘要
- 一句话:修复完全异步rollouter在部分rollout恢复时因drain loop阻塞导致无法及时分发新样本的问题。
- 推荐动作:该PR值得精读,特别是
_processor_worker方法中的drain循环重写展示了如何优雅处理异步任务管理和中断信号。关注点包括:
1) asyncio.Event与Condition的选择权衡;
2) 部分rollout启用/禁用时的不同处理策略;
3) 锁粒度优化以避免性能瓶颈。
功能与动机
根据PR body描述,当启用partial_rollout时,一旦_processor_worker进入drain循环(while self.active_tasks: ...),就会一直阻塞直到所有进行中任务完成。这导致新样本无法及时分发到空闲副本,影响了训练效率。修复目标是让drain循环可中断,以便在参数同步时能提前恢复样本分发。
实现拆解
- 同步原语重构:在
_init_async_objects方法中,将原有的asyncio.Condition和其内部锁替换为独立的asyncio.Lock和asyncio.Event。lock保护共享状态(paused/active_tasks/staleness_samples等),_resume_event作为运行状态信号。
- 状态重置逻辑更新:在
reset_staleness方法中,移除condition.notify_all()调用,改为设置_resume_event.set()来唤醒drain循环。
- drain循环重写:在
_processor_worker方法中,根据partial_rollout配置采用不同策略:
- 启用时:创建
resume_future监听恢复事件,使用asyncio.wait同时等待任务完成和恢复信号,收到信号后立即退出循环。
- 禁用时:保持原有逻辑但修复锁竞争问题。
- 无测试配套:本次变更未包含测试文件修改,属于核心逻辑修复。
关键文件:
verl/experimental/fully_async_policy/fully_async_rollouter.py(模块 完全异步;类别 source;类型 core-logic;符号 _init_async_objects, reset_staleness, _processor_worker): 这是本次PR唯一修改的文件,包含了完全异步rollouter的核心逻辑变更,特别是drain循环的重写和同步原语重构。
关键符号:_init_async_objects, reset_staleness, _processor_worker
关键源码片段
verl/experimental/fully_async_policy/fully_async_rollouter.py
这是本次PR唯一修改的文件,包含了完全异步rollouter的核心逻辑变更,特别是drain循环的重写和同步原语重构。
def _init_async_objects(self):
# Initialize asyncio synchronization primitives.
# `lock` protects shared state: paused / active_tasks / staleness_samples / timing fields.
self.lock = asyncio.Lock()
# `_resume_event` signals that the rollouter is currently running (paused == False).
self._resume_event = asyncio.Event()
self._resume_event.set() # 初始设置为运行状态
async def reset_staleness(self):
"""
Reset staleness samples after parameter update.
Returns timing_raw dictionary for metrics.
"""
async with self.lock:
self.paused = False
# 唤醒 _processor_worker 中的 drain 循环,使其能提前退出并恢复向空闲副本提交新样本
# 而不是等待所有进行中的长尾任务完成
self._resume_event.set()
# 每次参数更新后,重置陈旧样本计数
self.staleness_samples = len(self.active_tasks) + await self.message_queue_client.get_queue_size()
# ... 后续计时和日志逻辑保持不变
return timing_raw
async def _processor_worker(self):
"""
Streaming worker coroutines, a sample is submitted for processing without waiting for batches
"""
partial_rollout_enabled = self.config.async_training.get("partial_rollout", False)
while True:
if self.paused or await self._should_pause_generation():
print("[FullyAsyncRollouter][Processor] Received pause signal, waiting for remaining tasks...")
async with self.lock:
self.paused = True
self._resume_event.clear() # 清除事件,准备等待恢复信号
resume_future = asyncio.ensure_future(self._resume_event.wait())
try:
# Drain 循环:等待 (a) 至少一个活跃任务完成,或 (b) 恢复信号(reset_staleness/monitor 设置 paused=False)
# 以提前中断 drain,从而能向空闲副本提交新样本。
# 在等待期间不持有锁,因此发布者可以并发获取锁来更新 paused/staleness_samples。
while self.active_tasks and not resume_future.done():
wait_set = set(self.active_tasks) | {resume_future}
done, _pending = await asyncio.wait(wait_set, return_when=asyncio.FIRST_COMPLETED)
actual_done = done - {resume_future}
if actual_done:
async with self.lock:
for task in actual_done:
self.active_tasks.discard(task)
await task # 处理已完成任务
if resume_future in done:
print("[FullyAsyncRollouter][Processor] Resume event fired, breaking drain loop")
break
finally:
resume_future.cancel() # 确保清理 future
# ... 正常处理逻辑
评论区精华
gemini-code-assist[bot]在review中指出了几个关键问题:
- 锁竞争问题:在partial_rollout禁用时的else块中,
self.lock在await asyncio.wait期间被持有,会阻塞reset_staleness等方法的调用,导致训练管道性能停滞。
- 集合管理不一致:else块中直接替换
self.active_tasks集合引用,可能导致内存泄漏和任务丢失。
- 错误处理缺陷:对已完成任务的
await操作可能抛出异常。
建议采用统一的循环结构来提升健壮性和性能。ArronHZG随后批准了PR,表明问题已得到解决或接受。
- drain循环实现中的锁竞争和集合管理问题 (design): 建议采用统一的循环结构来提升健壮性和性能,ArronHZG的批准表明问题已解决或接受。
风险与影响
- 风险:
- 并发逻辑风险:新的drain循环逻辑涉及复杂的异步等待和锁管理,如果实现有误可能导致死锁或竞态条件。
- 向后兼容性:同步原语从Condition改为Lock+Event,可能影响依赖原有Condition通知机制的其他代码路径。
- 部分rollout配置依赖:修复仅针对
partial_rollout启用时生效,如果配置不一致可能无法解决原始问题。
- 测试覆盖不足:变更未包含测试更新,回归风险较高。
- 影响:
- 性能提升:解决了drain循环阻塞问题,使空闲副本能更快接收新样本,提高完全异步训练的资源利用率。
- 用户体验:训练流程更流畅,减少因等待长尾任务完成造成的人为停顿。
- 系统影响:仅影响完全异步rollouter模块,不涉及其他训练器或rollout后端。
- 团队影响:需要确保所有使用partial_rollout的团队了解此修复,并验证其训练脚本的兼容性。
- 风险标记:并发逻辑变更, 缺少测试覆盖, 配置依赖
关联脉络
- PR #6069 [fully_async] fix: fix rollouter/idle compute in async-mode: 同样修改了fully_async_rollouter.py文件,修复了空闲比计算逻辑,属于同一模块的连续修复。
- PR #6052 [fully_async] fix: avoid blocking ray.get inside async actor methods: 都涉及完全异步训练器的并发问题修复,关注异步环境下的阻塞和性能优化。
- PR #6046 [fully_async] fix: preserve per-iteration routed_experts on partial rollout resume: 都处理partial_rollout恢复时的逻辑问题,本PR修复样本分发,6046修复专家路由。
参与讨论