跳转至

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 间零拷贝传输,辅以共享内存作为兼容回退。