Prhub

#6491 [trainer] fix: gracefully shutdown trainer with TransferQueue

原始 PR 作者 wuxibin89 合并时间 2026-05-27 11:38 文件变更 1 提交数 1 评论 1 代码增减 +18 / -5

执行摘要

优雅关闭 trainer 的 ReplayBuffer 轮询线程

在分布式训练场景中,TransferQueue 的 kv_list 轮询线程在 trainer 退出时无法优雅停止,可能导致进程残留或异常退出。PR 旨在通过类似 grpc 的优雅关闭机制,确保背景线程在训练结束时有序终止。

建议尽快合并。该修复解决了训练终止时的资源泄漏和崩溃问题,对生产环境有积极意义。但 review 中提出的初始化失败场景仍需后续跟进修复,建议开一个新 issue 跟踪。

讨论亮点

Review 中 gemini-code-assist[bot] 指出一个临界 bug:当 PPOTrainer 初始化失败时(例如 _init_tokenizer_init_dataloader 抛异常),trainer 变量为 None,导致 trainer.replay_buffer.close() 不会执行,ReplayBuffer 后台线程泄露。随后 tq.close() 时该线程访问已关闭的队列会崩溃,掩盖原始错误。建议在 PPOTrainer.__init__ 内用 try...except 包裹初始化步骤,确保出错时也能关闭 replay buffer。

实现拆解

  1. ReplayBuffer.__init__ 中新增 self._stop_eventthreading.Event 对象),用于通知轮询线程停止。
  2. 修改 _poll_from_transfer_queue 方法:将 while True 改为 while not self._stop_event.is_set(),并将 time.sleep 替换为 self._stop_event.wait(self.poll_interval),使线程在等待期间可被事件唤醒。异常处理中增加 if not self._stop_event.is_set() 检查,避免在关闭过程中误判错误。
  3. 新增 close 方法:设置 _stop_event,然后 join 轮询线程(带超时),若超时则记录警告。
  4. run 方法的 finally 块中:先检查 trainer 是否为 None,再调用 trainer.replay_buffer.close(),确保无论 trainer 初始化是否成功,replay buffer 都能被关闭。
文件 模块 状态 重要度
verl/trainer/main_ppo_sync.py 训练器 modified 7.21

关键符号

close _poll_from_transfer_queue

关键源码片段

verl/trainer/main_ppo_sync.py core-logic

所有变更集中于此文件,包括 ReplayBuffer 的优雅关闭机制和 trainer run 方法中 finally 块的保护判断。

# verl/trainer/main_ppo_sync.py
# ReplayBuffer 类的关键变更片段
class ReplayBuffer:
    def __init__(self, poll_interval: float = 1.0):
        self.partitions: dict[str, dict[str, dict]] = defaultdict(dict)
        self.poll_interval = poll_interval
        self.lock = threading.Lock()
        # 新增:使用 Event 通知线程停止,替代无限循环
        self._stop_event = threading.Event()
        self.poll_thread = threading.Thread(
            target=self._poll_from_transfer_queue, daemon=True
        )
        self.poll_thread.start()
​
    def _poll_from_transfer_queue(self):
        """周期轮询 TransferQueue,支持优雅停止"""
        try:
            # 修改前 : while True:
            while not self._stop_event.is_set():
                data = tq.kv_list()
                if data is not None:
                    for partition_id, items in data.items():
                        self.add(partition_id, items)
                # 修改前 : time.sleep(self.poll_interval)
                # 使用 wait 替代 sleep,可被事件立即唤醒
                self._stop_event.wait(self.poll_interval)
        except Exception as e:
            # 仅当未收到停止信号时才视作错误
            if not self._stop_event.is_set():
                logger.error(f"Error in _poll_from_transfer_queue: {e}")
                os._exit(1)
​
    def close(self):
        """停止后台轮询线程"""
        if not self.poll_thread.is_alive():
            return
        self._stop_event.set() # 通知线程退出
        self.poll_thread.join(timeout=self.poll_interval + 1.0)
        if self.poll_thread.is_alive():
            logger.warning(
                "ReplayBuffer poll thread did not stop within timeout"
            )# run 方法中的 finally 块(部分)
    def run(self, config):
        # ... 初始化 ...
        trainer = None
        try:
            # ... 添加 worker 和初始化 trainer ...
            trainer.init_workers()
            trainer.fit()
        finally:
            # 新增保护:避免 trainer 初始化失败时访问 None
            if trainer:
                trainer.replay_buffer.close()
            tq.close()

评论区精华

trainer 初始化失败导致 replay buffer 泄漏 正确性

gemini-code-assist[bot] 指出,若 PPOTrainer.__init__ 抛出异常,trainer 为 None,replay_buffer.close() 不会执行,导致后台线程泄漏。建议在 __init__ 内用 try...except 包裹初始化并确保 close。

结论:本 PR 通过 `if trainer:` 部分缓解了问题,但未彻底解决 __init__ 内部异常场景。reviewer 建议后续跟进。 · unresolved

风险与影响

  1. 线程安全:使用 threading.Event 是标准模式,风险低。
  2. 超时逻辑join(timeout=poll_interval+1.0) 在极慢环境下可能不够,但已记录警告。
  3. review 指出的初始化失败问题:本 PR 仅修复了 finallyif trainer: 判断,但若 PPOTrainer.__init__ 内部抛出异常,replay buffer 仍可能泄漏。需要进一步将 close() 调用提前到 PPOTrainer.__init__ 的异常处理中。

直接改动仅一个文件(verl/trainer/main_ppo_sync.py),影响所有使用 PPO trainer 的训练流程,尤其是那些依赖 TransferQueue 进行跨进程通信的场景。优雅关闭可避免进程残留和资源泄漏,提升训练稳定性。

review 指出初始化失败时 replay buffer 泄漏未完全修复

关联 Issue

未识别关联 Issue

当前没有检测到明确关联的 Issue 链接,后续同步到相关引用后会出现在这里。

完整报告

参与讨论