跳转至

distributed_utils.py — 提供分布式通信工具

文件路径: verl/experimental/one_step_off_policy/distributed_utils.py

文件概述

提供分布式通信工具,主要是修补(monkey-patch)vLLM 的 StatelessProcessGroup.create 方法以支持 IPv6 和更灵活的网络配置。同时封装了 NCCL 通信组的初始化。

关键代码讲解

1. StatelessProcessGroup.create 修补

vLLM 的 StatelessProcessGroup 允许创建不污染 PyTorch 全局状态的进程组。原始实现不支持 IPv6,这里进行了增强:

@staticmethod
def create(host, port, rank, world_size, data_expiration_seconds=3600, store_timeout=300):
    # 自动检测地址类型
    try:
        ipaddress.IPv6Address(host.strip("[]"))
        address_family = socket.AF_INET6
    except (ipaddress.AddressValueError, ValueError):
        address_family = socket.AF_INET

    # 创建 socket
    if rank == 0:
        listen_socket = socket.socket(address_family, socket.SOCK_STREAM)
        if address_family == socket.AF_INET6:
            listen_socket.setsockopt(socket.IPPROTO_IPV6, socket.IPV6_V6ONLY, 1)
        listen_socket.bind((host.strip("[]"), port))
        listen_socket.listen()

    # 创建 TCPStore
    store = TCPStore(
        host_name=host, port=port, world_size=world_size,
        is_master=(rank == 0), timeout=timedelta(seconds=store_timeout),
        use_libuv=False, master_listen_fd=listen_fd,
    )
    return StatelessProcessGroup(rank=rank, world_size=world_size, store=store, ...)

# 应用修补
vllm.distributed.utils.StatelessProcessGroup.create = create

2. vllm_stateless_init_process_group - NCCL 初始化

def vllm_stateless_init_process_group(master_address, master_port, rank, world_size, device):
    """
    创建 StatelessProcessGroup 并初始化 NCCL 数据平面通信。
    用于训练进程和 vLLM 推理进程之间的权重同步。
    """
    if is_npu_available:
        from vllm_ascend.distributed.device_communicators.pyhccl import PyHcclCommunicator as PyNcclCommunicator
    else:
        from vllm.distributed.device_communicators.pynccl import PyNcclCommunicator

    pg = StatelessProcessGroup.create(host=master_address, port=master_port, rank=rank, world_size=world_size)
    pynccl = PyNcclCommunicator(pg, device=device)
    return pynccl

为什么需要 StatelessProcessGroup?

在 verl 中,训练进程和 vLLM 推理进程各自有自己的 torch.distributed 进程组。如果要在两者之间传输权重,需要建立额外的通信组。标准的 torch.distributed.init_process_group 是全局的,无法同时存在多个不同的组。StatelessProcessGroup 解决了这个问题——它创建独立的、不影响全局状态的进程组。

核心类/函数列表

名称 类型 说明
create() 静态方法 增强版的 StatelessProcessGroup 创建(支持 IPv6)
vllm_stateless_init_process_group() 函数 初始化无状态 NCCL 通信组

与其他模块的关系

  • 被权重同步模块使用(训练 GPU -> 推理 GPU 的参数传输)
  • 依赖 vLLM 的分布式通信基础设施
  • 支持 NPU(华为昇腾)设备

小结

distributed_utils.py 解决了"如何在已有分布式训练环境中建立额外通信通道"的问题。通过修补 vLLM 的 StatelessProcessGroup 增加了 IPv6 支持,并封装了 NCCL 通信组初始化。这是实现训练-推理分离架构中权重同步的基础设施。