bucketed_weight_transfer.py — 分桶权重传输¶
文件概述¶
实现基于 ZMQ + CUDA IPC(或共享内存回退)的分桶权重传输机制(约 302 行)。这是训练引擎向推理引擎同步权重的核心组件。
设计思想¶
模型权重可能非常大(几十 GB),无法一次性传输。分桶传输将权重分成固定大小的桶(默认 512MB),逐桶传输:
训练侧 (Sender) 推理侧 (Receiver)
│ │
│ 1. 初始化 ZMQ 连接 │
├────────────────────────────────────→│
│ │
│ 2. 发送 CUDA IPC handle │
├════════════════════════════════════→│ (重建 GPU buffer)
│ │
│ 3. 填充桶 → 发送元数据 │
├────────────────────────────────────→│ (从 buffer 读取权重)
│ │ (加载到模型)
│ 4. 填充桶 → 发送元数据 │
├────────────────────────────────────→│ (继续加载)
│ │
│ 5. 最后一桶 (is_last=True) │
├────────────────────────────────────→│ (完成)
│ │
核心类¶
BucketedWeightSender¶
class BucketedWeightSender:
"""发送端:打包权重到缓冲区,通过 ZMQ 发送"""
def __init__(self, zmq_handle, bucket_size_mb=512, use_shm=False):
self.bucket_size = int(bucket_size_mb) << 20 # MB → bytes
async def async_send_weights(self, weights):
"""异步发送权重
流程:
1. 初始化 ZMQ REQ socket
2. 创建 GPU buffer(或共享内存)
3. 发送 CUDA IPC handle 给接收端
4. 逐个权重填入 buffer,满了就发送
5. 发送最后一个桶
"""
for name, weight in weights:
if offset + weight.nbytes > self.bucket_size:
# 桶满了,发送并重置
socket.send_pyobj({"bucket_meta": bucket_meta, "is_last": False})
socket.recv() # 等待确认
offset = 0
# 将权重复制到 buffer
buffer[offset:offset+weight.nbytes].copy_(weight.view(torch.uint8))
offset += weight.nbytes
BucketedWeightReceiver¶
class BucketedWeightReceiver:
"""接收端:从 ZMQ 接收权重并回调处理"""
def receive_weights(self, on_bucket_received):
"""同步接收权重
流程:
1. 初始化 ZMQ REP socket
2. 接收 CUDA IPC handle,重建 buffer
3. 循环接收桶,每个桶解析出权重列表
4. 调用 on_bucket_received 回调处理
"""
while True:
metadata = socket.recv_pyobj()
weights = []
for name, meta in metadata["bucket_meta"].items():
tensor = buffer[offset:offset+size].view(dtype).view(shape)
if not use_shm:
tensor = tensor.clone() # CUDA IPC 需要 clone
weights.append((name, tensor))
on_bucket_received(weights)
if metadata["is_last"]:
break
CUDA IPC vs 共享内存¶
| 方式 | CUDA IPC | 共享内存 |
|---|---|---|
| 传输路径 | GPU → GPU(零拷贝) | CPU → GPU |
| 性能 | 高 | 较低 |
| 适用设备 | NVIDIA GPU | NPU/不支持 IPC 的设备 |
| 额外拷贝 | 需要 clone() 释放 IPC 内存 | 不需要 clone |
与其他模块的关系¶
BucketedWeightSender被vllm_rollout.py的ServerAdapter.update_weights()使用BucketedWeightReceiver被utils.py的vLLMColocateWorkerExtension.update_weights_from_ipc()使用- 通过 ZMQ IPC socket 通信(路径包含 GPU UUID)
小结¶
分桶权重传输是 Hybrid Engine 中训练→推理权重同步的核心机制,通过 CUDA IPC 实现 GPU 间零拷贝传输,辅以共享内存作为兼容回退。