Prhub

#7184 [rollout, vllm, hardware]fix: add IPC check before invoking ipc_collect in BucketedWeightReceiver

原始 PR 作者 kahlun 合并时间 2026-07-29 22:49 文件变更 1 提交数 1 评论 4 代码增减 +5 / -3

执行摘要

防止不支持 IPC 的设备上运行时崩溃

修复 bucketed weight receiver 在部分设备不支持 IPC 时调用 ipc_collect 导致运行时错误的问题。PR body 明确指出 'fix the runtime error when some device did not support IPC during bucketed weight receiver',并引用 vllm_rollout.py 中已有的类似 IPC 检查作为参考。

值得合并的快速修正。改动极少,但修复了真实运行时错误。可作为类似 IPC 调用的参考模式。

讨论亮点

审查者 wuxibin89 指出 'Sender side need guard as well.',作者 kahlun 随即在 BucketedWeightSender._cleanup 中增加了相同的 IPC 守卫,并合并为一次提交。这是一次有效且及时的审查反馈,避免了发送侧未修复导致的潜在问题。

实现拆解

  1. 导入 IPC 检测函数:在文件头部引入 is_support_ipc 函数(来自 verl.utils.device),用于运行时判断当前设备是否支持 IPC 操作。
  2. Receiver 侧清理加守卫:在 BucketedWeightReceiver._cleanup() 中,将直接调用 get_torch_device().ipc_collect() 改为先检查 is_support_ipc(),仅在返回 True 时才执行 ipc_collect()
  3. Sender 侧清理加守卫:根据审查意见,作者在 BucketedWeightSender._cleanup() 中也应用了相同的守卫,确保两端 IPC 清理路径行为一致。
  4. 测试说明:无新增测试文件,但作者在 A100 上运行了 tests/utils/test_bucketed_weight_transfer.py,所有 8 个 CUDA/IPC 测试均通过(5 个 SHM 测试因仅适用于 NPU 而跳过)。
文件 模块 状态 重要度
verl/workers/rollout/vllm_rollout/bucketed_weight_transfer.py 权重传输 modified 6.41

关键符号

_cleanup

关键源码片段

verl/workers/rollout/vllm_rollout/bucketed_weight_transfer.py dependency-wiring

所有变更均在此文件:导入 is_support_ipc,并在 BucketedWeightReceiver 和 BucketedWeightSender 的 _cleanup 方法中增加 IPC 守卫。

# verl/workers/rollout/vllm_rollout/bucketed_weight_transfer.pyimport gc
import logging
import os
from multiprocessing import shared_memory
from typing import Callable, TypedDictimport torch
import zmq
from torch.multiprocessing.reductions import reduce_tensorfrom verl.utils.device import get_device_id, get_device_name, get_torch_device, is_support_ipc# ... ( 中间省略 ) ...def _cleanup(self):
    """清理 socket、共享内存,并释放 GPU 缓存。"""
    if self.socket is not None:
        self.socket.close()
        self.socket = None
    if self.zmq_handle.startswith("ipc://"):
        ipc_path = self.zmq_handle[len("ipc://"):]
        try:
            os.remove(ipc_path)
        except OSError:
            pass
    del self.buffer
    self.buffer = None
    if self.shm is not None:
        self.shm.close()
        self.shm.unlink()
        del self.shm
        self.shm = None
    gc.collect()
    # 只有设备支持 IPC 时才调用 ipc_collect,否则跳过以避免运行时错误
    if is_support_ipc():
        get_torch_device().ipc_collect()
    get_torch_device().empty_cache()# 在 BucketedWeightReceiver._cleanup 中也应用了完全相同的守卫

评论区精华

Sender 侧同样需要 IPC 守卫 正确性

审查者 wuxibin89 评论 'Sender side need guard as well.',指出 BucketedWeightSender._cleanup 中也应该添加检查。

结论:作者 kahlun 在 Sender 侧增加了相同的 IPC 守卫,并合并为一次提交。 · 已解决

风险与影响

风险极低:变更仅限于两个 _cleanup 方法中增加条件判断,逻辑简单,未引入新的依赖或状态。但若 is_support_ipc 函数本身有 bug(比如误判或异常),可能导致清理路径意外跳过,但该函数在 vllm_rollout.py 中已有使用。

影响范围限定于 bucketed weight transfer 的清理路径,对正常传输流程无影响。用户侧仅限在非 CUDA/IPC 环境下运行的设备(如 NPU)受益,避免崩溃。对团队而言,这是一次清晰的防御性编程修正。

缺少测试覆盖 依赖新导入函数稳定性

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论