vllm-omni 中的 MoriTransferEngineConnector:基于 Mori IOEngine 的零拷贝 RDMA 阶段间传输实战指南
2026/9/17 3:22:39 网站建设 项目流程

vllm-omni 中的 MoriTransferEngineConnector:基于 Mori IOEngine 的零拷贝 RDMA 阶段间传输实战指南

【免费下载链接】vllm-omniA framework for efficient model inference with omni-modality models项目地址: https://gitcode.com/GitHub_Trending/vl/vllm-omni

本文以 vllm-omni 的 MoriTransferEngineConnector 设计文档 为核心,讲解如何在多阶段流水线(如 Thinker → Talker → Code2Wav)中通过 AMD 生态的 Mori 传输引擎实现 GPU 间零拷贝 RDMA 数据传输。你将掌握该连接器的适用场景、ZMQ 控制面与 RDMA 数据面的协同机制、deploy YAML 中的完整参数语义,以及基于 AMD MI300X 的 Qwen3-Omni-MoE 单节点(intra-node)可运行配置。

一、何时使用 MoriTransferEngineConnector

MoriTransferEngineConnector 是 vllm-omni 分布式连接器(OmniConnector)体系中的一员,位于vllm_omni/distributed/omni_connectors/connectors/mori_transfer_engine_connector.py,与 SharedMemoryConnector、MooncakeTransferEngineConnector、NixlConnector 等并列,通过 工厂注册机制 以MoriTransferEngineConnector名称创建。

适用边界

  • 当前仅支持单节点(intra-node)部署:在同一个节点内部,借助 AMD Infinity Fabric(XGMI)或本机 RDMA 网卡完成 GPU 间的数据搬运;
  • 跨节点(inter-node)支持暂缺:设计文档明确指出,跨节点支持会在未来的重构中恢复(文档引用了 issue #1742,inter-node 支持在重构中被移除)。

因此,若你的流水线分布在多台物理机上,当前应选择 MooncakeTransferEngineConnector 或 NixlConnector 等方案;若全部 stage 都在同一节点,且希望获得比共享内存(SHM)更低的 stage 切换延迟,则 Mori 是一个值得评估的选项。

与其他连接器的定位对比

从 disaggregated_inference.md 的连接器选型表中可以看到各方案的定位差异:

场景推荐连接器说明
单节点SharedMemoryConnector未显式配置时的默认连接器
多节点(Mooncake Store)MooncakeStoreConnectorTCP 传输,依赖 Mooncake Master 与 metadata 服务
多节点(Mooncake RDMA)MooncakeTransferEngineConnector带托管内存池的 RDMA/TCP 直传
多节点(Mori RDMA)MoriTransferEngineConnector经 Mori IOEngine 的 RDMA 直传
多节点(NIXL)NixlConnector基于 UCX 的 READ 传输,覆盖 Intel XPU
多节点(Yuanrong)YuanrongConnector依赖 Yuanrong Datasystem 与 etcd
Ascend NPU P2PYuanrongTransferEngineConnector直接使用 Yuanrong TransferEngine,需配置 NPU 设备 IPv4 与memory_pool_device: "npu"

二、工作原理:控制面与数据面分离

MoriTransferEngineConnector 的核心思路是用 ZMQ 做控制面、用 Mori 的 RDMA 做数据面,两者完全解耦:

  • 数据面(Data Plane):调用 Mori 的IOEngine/MemoryDescAPI,进行零拷贝 RDMA 传输。数据不经过 CPU 内存中转,直接从源 GPU 内存写入目标 GPU 内存;
  • 控制面(Control Plane):使用 ZMQ 完成拉取请求(pull-request)握手与异步完成通知。

从 源码实现 的模块文档可以提炼出三个关键设计点:

  1. 远程引擎必须先注册:任何数据传输开始前,必须通过register_remote_engine()注册对端(peer)的EngineDesc
  2. 内存区域用MemoryDesc描述:内存区域以可序列化的MemoryDesc对象(通过 pack/unpack 传输)而非裸虚拟地址标识,天然支持跨进程传递;
  3. 传输经batch_write()派发:数据传输通过IOEngine.batch_write()发起,并借助TransferStatus对象进行异步跟踪。

ZMQ 握手协议

控制面定义了两种请求类型(源码定义):

  • Pull 请求(MoriPullRequest):接收方(receiver)向发送方(sender)发送一个 msgspec 编码的结构体,包含自身的EngineDesc与内存池MemoryDesc,请求发送方将数据直接 RDMA 写入接收方内存池的指定偏移(dst_offset);
  • Query 请求(QueryRequest):接收方向发送方发送QUERY_INFO前缀 +QueryRequest,查询指定request_id的数据是否已就绪——这是无 metadata 的 get 路径(见下文 get 流程)。

ZMQ 消息中还会用到四个固定字节标记:trans_done(传输成功)、trans_error(传输失败)、query_info(查询请求前缀)、info_not_found(数据不存在)。

底层传输:backend 可选

Mori 的IOEngine.batch_write()是 backend 无关的,构造时通过create_backend()选择具体传输通道(源码):

  • rdma:基于网卡的 RDMA(RoCE / InfiniBand,支持 GDR 直访),需要配置device_name(如mlx5_0)或通过环境变量MORI_RDMA_DEVICES指定;
  • xgmi:使用 AMD Infinity Fabric 的 GPU-to-GPU 直连链路,因此必须使用 CUDA 内存池(见下文参数校验)。

无论选择哪种 backend,数据面路径与 ZMQ 握手逻辑完全一致,只是内存池的物理位置与网卡/互联设备的差异。

三、安装与依赖

MoriTransferEngineConnector 对 Mori 库采用惰性导入:源码在try块中导入mori.cpp.TransferStatusmori.io(IOEngine、EngineDesc、MemoryDesc、RdmaBackendConfig、XgmiBackendConfig 等),导入失败时IOEngine = None,并在构造函数中抛出明确的ImportError(源码)。因此需要先安装 Mori(官方为pip install mori)。

以 AMD MI300X 单节点部署为例,参考 qwen3_omni_moe_mori_intranode.yaml 的依赖说明,运行环境通常需要:

  • moripip install mori
  • torch-rocm >= 2.9
  • 带有 #2383 deploy schema 的 vllm-omni
  • 通过/dev/infiniband暴露的 ibverbs / RDMA 网卡(XGMI backend 下可不依赖,但 RDMA backend 必须有)

四、配置:deploy YAML 完整参数语义

Mori 连接器通过新的 deploy-config schema 配置(见 stage_configs.md)。配置方式为:在 deploy YAML 顶层定义connectors,并在各 stage 的input_connectors/output_connectors中按名称引用。

最小配置示例(源自设计文档)

connectors: mori_connector: name: MoriTransferEngineConnector extra: host: "auto" zmq_port: 50051 device_name: "" memory_pool_size: 536870912 memory_pool_device: "cuda" stages: - stage_id: 0 output_connectors: to_stage_1: mori_connector - stage_id: 1 input_connectors: from_stage_0: mori_connector

配置经过 ConnectorSpec 数据类解析:name用于在工厂注册表中定位构造函数,extra字典原样传给MoriTransferEngineConnector.__init__(factory.create_connector)。

参数详解

结合 连接器源码 的解析逻辑,各参数含义、默认值与约束如下:

参数默认值说明与约束
host"127.0.0.1"本机 RDMA IP;设为"auto"时自动探测(通过 UDP socket 连接外部地址获取本机 IP,失败时回退到socket.gethostbyname,再失败回退127.0.0.1
zmq_port50051ZMQ 控制面基础端口,sender 在此端口上绑定 ROUTER socket
backend_type"rdma"可选"rdma"(NIC 上的 RoCE/IB,GDR 能力)或"xgmi"(AMD Infinity Fabric GPU 直连);非法值直接抛ValueError
device_name""RDMA 设备名(如"mlx5_0");仅在rdmabackend 下生效,空串时读取环境变量MORI_RDMA_DEVICESxgmi下设置该值会被忽略并告警
memory_pool_size1073741824(1 GiB)RDMA 内存池大小(字节);设计文档示例用536870912(512 MB),与 Qwen2.5-Omni 的 Mori 部署保持一致
memory_pool_device"cpu""cpu"表示 pinned 内存(torch.empty(...).pin_memory());"cuda"表示 GPU 内存(GPUDirect / XGMI RDMA 必需)
role"sender""sender"绑定 ZMQ 监听并接受put()"receiver"不监听,仅通过查询上游 sender 的sender_host/sender_zmq_port执行get()。非法值抛ValueError
sender_host/sender_zmq_portNonereceiver 侧用于定位上游 sender 的 ZMQ 端点;为None"auto"时等待update_sender_info()注入
xgmi_num_streams/xgmi_num_events64/64xgmibackend 生效,配置XgmiBackendConfig的流与事件数量
qp_per_transfer1rdmabackend 生效,每次传输使用的队列对数
post_batch_size-1rdmabackend 生效,RdmaBackendConfig的批处理大小
num_worker_threads1rdmabackend 生效,RDMA worker 线程数

关键约束(源码级校验)

  1. backend 与内存池必须匹配backend_type='xgmi'memory_pool_device='cpu'时直接抛ValueError——XGMI 是 GPU 到 GPU 的互联,无法寻址 CPU 内存(源码);
  2. 角色由role显式指定:sender 绑定 ZMQ 失败(端口占用、地址不可用、权限不足等)属于致命错误,会快速失败并在__init__中传播异常,不存在静默回退为 receiver 的逻辑;
  3. 端点注入:receiver 可通过update_sender_info(sender_host, sender_zmq_port)在运行期注入 sender 的 ZMQ 端点(源码)。

五、实战示例:Qwen3-Omni-MoE 在 AMD MI300X 上的 intra-node 部署

仓库提供了开箱即用的示例 qwen3_omni_moe_mori_intranode.yaml,面向单个 AMD Instinct MI300X OAM 节点(8 卡、192GB HBM3、gfx942)。该配置与 CUDA 侧的qwen3_omni_moe.yaml(走 SharedMemoryConnector)保持相同的流水线拓扑与采样参数,仅将连接器替换为MoriTransferEngineConnector,并以backend_type: xgmi走 AMD Infinity Fabric GPU-to-GPU 直连路径。

启动命令

vllm-omni serve Qwen/Qwen3-Omni-30B-A3B-Instruct --omni --log-stats \ --deploy-config vllm_omni/deploy/qwen3_omni_moe_mori_intranode.yaml

完整 YAML 剖析

async_chunk: true connectors: mori_connector: name: MoriTransferEngineConnector extra: host: "auto" zmq_port: 50051 backend_type: "xgmi" # AMD Infinity Fabric GPU-to-GPU 直连 device_name: "" # 留空以尊重 $MORI_RDMA_DEVICES memory_pool_size: 536870912 # 512 MB,与 Qwen2.5-Omni Mori 部署一致 memory_pool_device: "cuda" # XGMI RDMA 传输要求 GPU 内存 xgmi_num_streams: 64 xgmi_num_events: 64 qp_per_transfer: 1 num_worker_threads: 1 post_batch_size: -1 # chunk 传输参数,供 talker2code2wav_async_chunk 消费: # 每次 connector.put 携带 25 帧解码后的 codec 帧(约 625Hz codec 速率下约 40ms), # 并带 25 帧左上下文以满足 code2wav 的感受野。 # 与 SHM 部署保持一致,保证切换 backend 后精度可比。 codec_chunk_frames: 25 codec_left_context_frames: 25 stages: - stage_id: 0 gpu_memory_utilization: 0.9 enforce_eager: true # ROCm 覆盖:规避 flashinfer autotune 路径 devices: "0" output_connectors: to_stage_1: mori_connector default_sampling_params: temperature: 0.4 top_p: 0.9 top_k: 1 max_tokens: 2048 seed: 42 repetition_penalty: 1.05 - stage_id: 1 gpu_memory_utilization: 0.9 enforce_eager: true devices: "1" input_connectors: from_stage_0: mori_connector output_connectors: to_stage_2: mori_connector default_sampling_params: temperature: 0.9 top_k: 50 max_tokens: 4096 seed: 42 repetition_penalty: 1.05 - stage_id: 2 gpu_memory_utilization: 0.3 max_num_seqs: 1 enforce_eager: true async_scheduling: false # Codec prefill 长度(Q * num_frames)超过默认 32k; # 与 qwen3_omni_moe.yaml 的 code2wav 尺寸保持一致。 max_num_batched_tokens: 51200 devices: "2" input_connectors: from_stage_1: mori_connector default_sampling_params: temperature: 0.0 top_p: 1.0 top_k: -1 max_tokens: 65536 seed: 42 repetition_penalty: 1.1

该配置的关键设计点

  1. 三阶段流水线:stage 0(Thinker)→ stage 1(Talker)→ stage 2(Code2Wav),每阶段独占一块 GPU(devices: "0"/"1"/"2");
  2. async_chunk: true:启用 chunk 级传输,接入OmniChunkTransferAdapter(chunk_transfer_adapter.py)。stage 间的 hidden-state 流与 codec 帧流以调度器拥有的后台 put/get 线程异步搬运,这正是测试 Mori 数据面的主要路径;
  3. codec_chunk_frames: 25/codec_left_context_frames: 25:控制每次connector.put()携带的 codec 帧数与左上下文帧数,参数与 SHM 部署保持一致,保证后端切换后精度可比;
  4. enforce_eager: true:ROCm 下的必要覆盖,用于规避 flashinfer autotune 路径;
  5. max_num_batched_tokens: 51200:因为 codec prefill 长度(Q × num_frames)超过默认 32k 上限。

为什么选择 Mori 而非 SHM

YAML 头部注释给出了明确的动机:在 CUDA/H100 上 Qwen3-Omni-MoE 通过 SharedMemoryConnector 传输(见qwen3_omni_moe.yaml);而在 MI300X 上 SHM 路径虽然可用,但XGMI RDMA 可以显著降低 stage 切换延迟,尤其针对两条 async-chunk 模式下的关键数据流:

  • thinker → talker 的 hidden-state 流;
  • talker → code2wav 的 codec 帧流。

这些数据流通过 XGMI 走 GPU-to-GPU 直传,绕过了 CPU/主机内存的中转。

已知限制

  • 仅支持 TP=1:连接器在tp_world_size > 1时会抛NotImplementedError(chunk 路径的端点推导尚未做到 rank 感知),配置注释明确标注了这一限制;
  • Qwen2.5-Omni + Mori 的 chunk 路径尚不可用:它需要thinker2talker_async_chunk/talker2code2wav_async_chunk输入处理器,而这些处理器目前尚不存在(上游 PR 原本针对的 orchestrator 级路径在 entrypoints → engine 重构中丢失,并随 #1742 移除)。

六、源码级深入:put/get 调用链

统一接口

所有 OmniConnector 都继承 OmniConnectorBase,实现四个抽象方法:put()get()cleanup()health()close()MoriTransferEngineConnector.supports_raw_data = True,意味着它能以原生方式处理bytes/torch.Tensor载荷而无需走序列化器。

put():发送方路径

put()(源码)的流程如下:

  1. 拒绝空载荷(空bytes/ 空torch.Tensor);
  2. 非原生类型(非ManagedBuffer/torch.Tensor/bytes)通过OmniSerializer.serialize()序列化,并标记is_fast_path = False
  3. 分配内存池中的一块缓冲(BufferAllocator 以 4096 字节对齐管理空闲块),把源数据拷贝进池中(跨设备时使用non_blocking拷贝并显式同步 CUDA 流);
  4. {put_key}@{from_stage}_{to_stage}的键格式将(offset, size, holder, ...)记录到_local_buffers,并返回(True, size, metadata),其中 metadata 携带source_hostsource_portdata_sizeis_fast_path
  5. 同一 put_key 的旧缓冲会被弹出并释放,防止重复键导致的内存泄漏。

get():接收方路径

get()(源码)支持两种模式:

  • 带 metadata:直接从 metadata 中取source_host/source_port/data_size/is_fast_path,跳过查询握手;
  • 无 metadata:先通过 ZMQ 向sender_host:sender_zmq_port发送QUERY_INFO前缀 +QueryRequest,收到info_not_found则返回None,收到QueryResponse则回填 metadata。

拿到 metadata 后,接收方构建MoriPullRequest(携带自身engine_desc_packed、内存池mem_desc_packeddst_offsetlength),发送给 sender;sender 的监听线程校验数据长度后调用engine.batch_write()发起 RDMA 写;接收方收到trans_done后:

  • is_fast_path = True时直接返回ManagedBuffer(零拷贝消费);
  • 否则反序列化载荷并释放缓冲。

超时策略为 30 秒基础超时 + 按数据量线性增加(每 100 MB 加 5 秒)。

sender 侧监听线程

sender 在role='sender'时启动 ZMQ ROUTER 监听线程(源码),配合 4 线程的ThreadPoolExecutor异步处理 pull 请求与 query 请求,并通过inproc通道唤醒监听循环发送响应。两个值得注意的健壮性设计:

  1. 长度校验防数据损坏pull.length != src_size时拒绝传输——长度偏小会静默截断载荷,偏大则因pool_mem_desc覆盖整个内存池而越界读取相邻的在途缓冲,因此宁可拒绝也不静默损坏数据;
  2. TTL 清理防内存泄漏:本地缓冲超过 300 秒(_BUFFER_TTL_SECONDS)未消费即被回收,防止接收方崩溃或超时后数据永久滞留(源码注释同时承认了 TTL 清理与在途 RDMA 传输存在竞态的极端场景,留待后续 PR 处理)。

运维接口

  • health():返回statushostpool_devicepool_sizeengine_key以及puts/gets/bytes_transferred/errors/timeouts五项指标,便于监控排查;
  • close():幂等的资源释放——停止监听线程与线程池、释放本地缓冲、关闭 ZMQ 上下文,并通过deregister_memory/deregister_remote_engine尽力反注册 Mori 资源后置空engine引用,交由 Python GC 触发 C++ 析构(Mori 的IOEngine不暴露显式 shutdown 接口)。

七、与 chunk 传输适配器的协作

async_chunk: true时,连接器由 OmniChunkTransferAdapter 驱动:调度器通过save_async()将多模态输出封装为 chunk 任务入队,后台 save 线程执行connector.put();下游 stage 通过load_async()注册请求,后台 recv 线程以connector.get()非阻塞轮询(timeout=0),请求在 chunk 到达期间进入WAITING_FOR_CHUNK状态,收到带finished标记的终止 chunk 后才结束。

这条路径意味着 Mori 连接器的put/get会被调度器拥有的后台线程并发调用,因此实现中大量使用_local_buffers_lock_registered_engines_lock等细粒度锁,并在_poll_single_request/cleanup_receiver之间建立了提交屏障(commit barrier),确保清理与在途 chunk 提交互不冲突。

八、总结

MoriTransferEngineConnector 为 vllm-omni 在 AMD 平台上提供了一条绕过主机内存的 GPU-to-GPU 数据传输通道:ZMQ 负责握手与状态查询,MoriIOEngine.batch_write()负责真正的零拷贝 RDMA 写入,配合 4096 字节对齐的托管内存池、300 秒 TTL 清理与严格的长度校验,构成了一个工程上相当完整的传输引擎连接器。当前它的成熟战场是单节点 + XGMI(AMD MI300X 上的 Qwen3-Omni-MoE 即为此而设计),TP=1 限制与 Qwen2.5-Omni chunk 路径缺失是需要提前确认的前提条件。如需进一步了解 deploy schema 的通用语法与其余连接器,可继续阅读 stage_configs.md 与 disaggregated_inference.md。

【免费下载链接】vllm-omniA framework for efficient model inference with omni-modality models项目地址: https://gitcode.com/GitHub_Trending/vl/vllm-omni

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询