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 通信组初始化。这是实现训练-推理分离架构中权重同步的基础设施。