ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Mooncake EP 与 Mooncake PG 实战指南:为容错 MoE 推理构建专家并行调度与动态成员通信后端

Mooncake EP 与 Mooncake PG 实战指南:为容错 MoE 推理构建专家并行调度与动态成员通信后端 人工智能大模型模型推理服务后端【免费下载链接】MooncakeMooncake is the serving platform for Kimi, a leading LLM service provided by Moonshot AI.项目地址https://gitcode.com/gh_mirrors/mo/Mooncake点击查看免费下载导读Mooncake 为 MoE 推理提供了两个紧密配合的组件Mooncake PG是一个基于torch.distributed的 ProcessGroup 后端注册mooncake加速器后端与mooncake-cpu后端实现集合通信、点对点通信以及动态成员elastic membership管理Mooncake EP则是面向低延迟 MoE 推理的专家并行Expert Paralleldispatch/combine 运行时遵循 DeepEP 低延迟编程模型并叠加了 rank 活跃度感知与 Mooncake 传输支持。读完本文你将掌握两个组件的构建方式、PG 的初始化与弹性恢复流程、EP Buffer 的 dispatch/combine 完整用法以及基于 NCCL 的ElasticBuffer传输自动选择机制并能在 PG/EP 故障排查指南 的辅助下定位常见问题。组件定位PG 负责成员关系EP 负责专家流量两者的典型集成模式是先初始化一个 Mooncake process group再基于该 group 构造 Mooncake EP 的Buffer。Process group 既用于常规集合通信也用于交换 EP 的启动元数据bootstrap metadata。Mooncake PG注册mooncake与mooncake-cpu两个torch.distributed后端实现 collective 与 point-to-point API并暴露动态成员辅助函数。Mooncake EP专家并行 dispatch/combine 运行时用于延迟敏感的 MoE 推理支持 rank 活跃度感知与 Mooncake 传输。更深层的设计细节可参见 Mooncake PG 设计文档 与 Mooncake EP 设计文档。构建与安装Mooncake EP 与 PG 包含在启用 CUDA 的 Mooncake wheel 中。从源码构建时需要显式启用 EP/PG 扩展cmake .. -DWITH_EPON make -j扩展是针对特定 PyTorch 版本编译的。在导入时mooncake.pg与mooncake.ep会加载带版本后缀的扩展模块使其与当前的torch.__version__匹配。如果当前 PyTorch 版本与已构建的扩展不匹配导入会失败并提示如Mooncake PG was not built against torch...的错误。从源码可见其实现机制python/mooncake/ep.py 首先将torch.__version__解析为形如_2_5_1的版本后缀然后分别导入mooncake._ep与mooncake.pgversion_suffix两个模块任一缺失都会抛出对应的ImportError。因此务必保证运行时使用的 Python 环境与构建时一致。NCCL 支持是可选opt-in的需要额外开启-DWITH_EPON -DUSE_CUDAON -DUSE_NCCL_DEVICEON该选项默认为OFF详见下文ElasticBuffer 默认 NCCL 后端。Mooncake PG 快速开始CUDA 后端import os import torch import torch.distributed as dist from mooncake import pg rank int(os.environ[RANK]) world_size int(os.environ[WORLD_SIZE]) local_rank int(os.environ.get(LOCAL_RANK, rank)) torch.cuda.set_device(local_rank) device torch.device(cuda, local_rank) dist.init_process_group( backendmooncake, rankrank, world_sizeworld_size, ) x torch.tensor([rank 1], dtypetorch.int32, devicedevice) dist.all_reduce(x, opdist.ReduceOp.SUM) print(frank{rank}, all_reduce{int(x.cpu())})使用常规的 PyTorch launcher 运行例如torchrun --nproc-per-node2 pg_quickstart.pyCPU 后端将 backend 换为mooncake-cpu即可dist.init_process_group( backendmooncake-cpu, rankrank, world_sizeworld_size, )对于使用默认故障处理方式的固定规模 grouppg_options是可选的。当需要预留额外 group 容量、以扩展成员身份加入extension join或选择非默认的故障处理模式时才需要传入MooncakeBackendOptions。选择网络设备要显式将 Mooncake 限制到指定的 NIC / HCA 设备列表需要在init_process_group()之前调用pg.set_device_filter(...)from mooncake import pg pg.set_device_filter([mlx5_1, mlx5_2])在测试与基准命令中同样的设置通常通过环境变量MOONCAKE_PGTEST_DEVICE_FILTERSmlx5_1,mlx5_2传递。这在多 RDMA 设备机器上尤其重要——如果后端自动选中了用于管理流量、位于不同 fabric 或对端不可达的设备可能出现连接建立失败、反复超时、意外回退到 fallback 或带宽异常等问题显式设置 filter 后问题消失即可判定为拓扑/设备选择问题而非集合通信正确性问题。Mooncake PG Torch API 参考MooncakeBackendOptionspg.MooncakeBackendOptions(max_group_size) pg.MooncakeBackendOptions(max_group_size, is_extension) pg.MooncakeBackendOptions( max_group_size, is_extension, auto_deactivate_on_failure, auto_sync_on_failure, ) # Explicit active-rank mirror overloads pg.MooncakeBackendOptions(active_ranks) pg.MooncakeBackendOptions(active_ranks, is_extension) pg.MooncakeBackendOptions(active_ranks, is_extension, max_group_size)参数说明max_group_sizegroup 内固定的槽位slot容量。它必须不小于初始声明的 group 规模且之后不能再增大。active_ranks可选的连续torch.int32存储作为已提交 PG 成员关系的镜像。其初始内容会被忽略按max_group_size尺寸分配可位于 CPU 或 GPU。is_extension对于通过join_group()进入既有 group 的替换或加入进程需要设为True。auto_deactivate_on_failure与auto_sync_on_failure选择自动式还是框架托管式故障处理。两者默认均为True自动同步auto-sync依赖自动停用auto-deactivation。实用工具函数函数用途备注pg.set_host_ip(host_ip)覆盖后端使用的 host IP。在init_process_group()之前调用。pg.set_device_filter(filters)限制 NIC/HCA 选择。在init_process_group()之前调用。pg.set_transfer_engine(engine)复用外部TransferEngine。该 engine 必须比所有 process group 存活得更久。pg.get_active_ranks(backend)返回后端的 active-rank 张量。用于 EP 的 fallback 与恢复路径。pg.get_num_synced_ranks(backend)返回本地可激活的 group 槽位数量。诊断辅助函数。pg.get_peer_state(backend, ranks)读取本地镜像的激活就绪状态。轻量、无通信的查询。pg.activate_ranks(backend, ranks)通过 Coordinator 提议激活。任意一个在线 rank 调用一次即可。pg.recover_ranks(backend, ranks)通过 Coordinator 提议激活。activate_ranks的兼容别名。pg.deactivate_ranks(backend, ranks)通过 Coordinator 提议停用。任意一个在线 rank 调用一次即可。pg.join_group(backend)确认激活就绪并保持阻塞直到激活真正发生。用于 scale-up、进程替换与原地 rejoin。pg.sync_after_failure(backend)上报当前链路观测、等待对账并应用最新 group 视图。当auto_sync_on_failureTrue时自动调用也可手动调用。支持的分布式操作Mooncake PG 实现了以下torch.distributedAPI。具体支持程度可能依赖设备类型、dtype、PyTorch 版本以及当前后端是mooncake还是mooncake-cpu生产使用前请在目标环境上运行 PG 测试验证。API 族示例备注Collectivesall_reduce、broadcast、all_gather、all_gather_into_tensor、reduce_scatter_tensor、all_to_all、barrier、reduce、gather、scatter活跃 rank 参与不活跃 rank 会被后端内部跳过。Async workdist.all_reduce(..., async_opTrue)等待返回的 work 对象然后按需同步设备流。P2Pisend、irecv、batch_isend_irecv单张量的 P2P 通过 Mooncake 后端 shim 路由。从架构设计文档看PG 的 P2P 采用接收方驱动的信用credit协议接收操作从接收池预留 chunk 并向发送方写CreditSlot发送方在发送池中暂存数据、通过 Transfer Engine 写入目的地址并回写AckSlot接收方再拷贝到用户缓冲控制槽位携带 group epoch 与序号配对的 header/footer token 防止半写槽位被消费epoch 校验会丢弃旧 group view 遗留的控制流量。集合通信则采用基于 Transfer Engine 的直接写direct-write设计每个 communicator 注册自己的发送/接收/同步缓冲详见 mooncake-backend-pg.md。弹性恢复协议Mooncake PG 将加入准备join preparation与成员激活membership activation分离加入或恢复的 rank 完成本地 warmup 后调用join_group()并等待已有 rank 可以轮询本地就绪状态然后发出激活提议Coordinator 负责校验并分发最终的成员关系。健康 rank 一侧from mooncake import pg dist.init_process_group( backendmooncake, rankrank, world_size2, pg_optionspg.MooncakeBackendOptions( 3, # max_group_size False, # is_extension ), ) backend dist.group.WORLD join_ranks [2] while not all(pg.get_peer_state(backend, join_ranks)): # Continue serving, back off, or poll according to your scheduler policy. pass pg.recover_ranks(backend, join_ranks)加入 rank 一侧from mooncake import pg dist.init_process_group( backendmooncake, rank2, world_size3, pg_optionspg.MooncakeBackendOptions( 3, # max_group_size True, # is_extension ), ) backend dist.group.WORLD # Collectives are local-only before join_group. Use this # window for framework-specific preparation, for example: # capture_cuda_graphs() # warm_up_model() pg.join_group(backend)关键语义get_peer_state()是本地尽力而为best-effort的就绪查询不是集合操作。容量必须在创始成员创建 group 时通过max_group_size预留。加入的注册会在该容量内追加不活跃槽位。加入的 rank 在join_group之前只表现为本地集体行为join 调用随后会阻塞直到 Coordinator 批准的激活提交完成。任意在线 rank 调用一次activate_ranks()或其别名recover_ranks()即足够冗余的等价调用是安全的。子 group 必须在健康进程与加入进程上按相同顺序创建遵循 PyTorchnew_group()的排序规则。从设计文档理解动态成员模型PG 采用带编号座位的餐桌心智模型桌子座位数在开席前固定即max_group_size容量就餐过程中座位不一定坐满——成员可以离席但座位不消失、其余人不换号之后该成员可回到同一座位或由新成员占用也可以占用预留但从未使用的座位。因此桌子有多少座位与哪些座位当前有人参与是两个独立问题前者固定、后者随时间变化。例如一个 group 预留 8 个槽位、最初只激活 4 个max_group_size 8 active_ranks [1, 1, 1, 1, 0, 0, 0, 0]rank 2 被停用后active_ranks [1, 1, 0, 1, 0, 0, 0, 0]其槽位仍被保留其他 rank 不会被重新编号。rank 4 随后也可在不填补该空洞的情况下变为活跃active_ranks [1, 1, 0, 1, 1, 0, 0, 0]需要特别注意的是group 的size并不是活跃成员数而是最高活跃 in-group rank 加一它作为 rank 索引的上界。例如active_ranks [1, 0, 1, 0]时size 3、max_group_size 4。按 rank 索引分配缓冲的用户代码必须覆盖至少size个条目包括该上界以下的空洞。此外还有进程级max_world_size与 group 级max_group_size两层容量概念前者约束整个 world 的 rank 数后者约束单个 group 的 rank 数。Coordinator 还维护进程级RankStateOffline/Synced/Healthy与 group 级GroupMemberStateNone/Inactive/AwaitingActivation/Active/Left。创始成员在 bootstrap 时直接成为Active加入成员的激活路径为Inactive → AwaitingActivation → Active。激活需满足完整条件group 就绪、每个新激活目标处于AwaitingActivation且Healthy且已发布 endpoint、并且结果活跃集合内任意两个 rank 之间双向连通。Coordinator 目前运行在全局 rank 0 上不提供高可用。故障处理模式auto_deactivate_on_failure决定失败驱动的成员变更由谁负责Mooncake PG 或上层框架auto_sync_on_failure决定失败的 collective/P2P 操作在完成前是否自动调用sync_after_failure。存在三种合法配置与一种非法组合PG 托管、同步默认auto_deactivate_on_failuretrueauto_sync_on_failuretrue。失败操作上报观测并自动进入sync_after_failureCoordinator 对账后停用不健康成员并返回新视图调用方在失败操作完成前应用该视图。这是最安全简单的模式但把控制面延迟放进了失败完成路径负向观测会打开对账窗口失败操作可能挂起数十秒等待跨 rank 对账。PG 托管、延迟同步auto_deactivate_on_failuretrueauto_sync_on_failurefalse。操作在数据面工作完成后即返回暴露local_success与failed_ranks_hint框架可先做其他工作、再在依赖新成员关系或恢复通信前调用sync_after_failure。框架托管auto_deactivate_on_failurefalseauto_sync_on_failurefalse。失败操作只返回本地证据由框架观察local_success/failed_ranks_hint决定移除哪些 rank再调用deactivate_ranks。非法组合auto_deactivate_on_failurefalseauto_sync_on_failuretrue在构造时会被拒绝因为只有 PG 同时拥有失败驱动的停用权自动同步才有意义。此限制只针对自动同步sync_after_failure在任何模式下都可手动调用。失败操作会向调用者返回两个结果local_success本次操作在调用方完成全部传输与否与failed_ranks_hint长度为max_group_size、按 in-group rank 索引的逐操作位图置位表示该 rank 出现传输失败——这是证据而非全局结论。负向证据会在 Coordinator 打开默认 30 秒的对账窗口需配置为大于默认 collective 超时窗口关闭后 Coordinator 推导双向连通的健康集合、更新RankState并分发结果启用 auto-deactivation 时还会把不健康活跃成员改为Inactive。Mooncake EP 快速开始Mooncake EP 从mooncake.mooncake_ep_buffer暴露Buffer使用一个 Mooncake process group 和按预期 dispatch 形状计算的工作区大小来初始化import torch import torch.distributed as dist from mooncake import pg from mooncake.mooncake_ep_buffer import Buffer # Assume dist.init_process_group(..., backendmooncake, ...) has completed. group dist.group.WORLD rank dist.get_rank(group) world_size dist.get_world_size(group) num_tokens 128 hidden 7168 num_experts 288 top_k 8 max_tokens_per_rank 128 x torch.randn(num_tokens, hidden, dtypetorch.bfloat16, devicecuda) scores torch.randn(num_tokens, num_experts, dtypetorch.float32, devicecuda) topk_idx torch.topk(scores, top_k, dim-1).indices topk_weights torch.softmax( torch.randn(num_tokens, top_k, dtypetorch.float32, devicecuda), dim-1 ) num_ep_buffer_bytes Buffer.get_ep_buffer_size_hint( max_tokens_per_rank, hidden, world_size, num_experts, ) buffer Buffer(group, num_ep_buffer_bytes) # EP-level rank-health tensor. Kernels may update it to 0 when timeout_us # detects a failed source rank. active_ranks torch.ones(world_size, dtypetorch.int32, devicecuda) recv_x, recv_count, handle, event, hook buffer.dispatch( x, topk_idx, active_ranks, num_max_dispatch_tokens_per_rankmax_tokens_per_rank, num_expertsnum_experts, timeout_us-1, use_fp8True, async_finishFalse, return_recv_hookFalse, ) event.current_stream_wait() # Run local experts on recv_x here. If use_fp8True, recv_x is a # (data, scales) tuple; dequantize or feed it into an FP8-aware expert kernel. expert_out run_local_experts(recv_x, recv_count) combined_x, event, hook buffer.combine( expert_out, topk_idx, topk_weights, active_ranks, timeout_us-1, handlehandle, ) event.current_stream_wait()EP 的整体数据流包含三个阶段Dispatch各 rank 将 token 隐状态发给拥有被选中专家的 rank接收方按本地专家打包 token→Expert compute每个 rank 在打包输入上运行本地专家→Combine专家输出按路由权重回传并归约到原始 token 所有者。Mooncake EP API 参考Buffer.get_ep_buffer_size_hint(...)Buffer.get_ep_buffer_size_hint( num_max_dispatch_tokens_per_rank: int, hidden: int, num_ranks: int, num_experts: int, ) - int返回 EP buffer 的工作区字节数。应使用单个 rank 单步可能 dispatch 的最大 token 数来调用低估该值可能导致 buffer 溢出或 dispatch 结果错误设计文档同样强调按峰值 dispatch 需求而非平均请求来配置。Buffer(group, num_ep_buffer_bytes0)为 Mooncake process group 创建 EP 运行时。构造函数通过 group 交换 RDMA 与 IPC 元数据在可用时初始化快速路径传输若快速路径不可用则回退到基于 PyTorch collectives 的 Python 实现。从 python/mooncake/mooncake_ep_buffer.py 的实现看connect()依次执行all_gather交换本地 RDMA memory-region 地址raddr与 keyrkey、all_to_all交换本地 QPN 与 LID、all_gather交换 GID 的 subnet prefix 与 interface ID最后用sync_ibgda_peers()同步对端若p2p_enabled()则再all_gather交换 CUDA IPC handle 并调用sync_nvlink_ipc_handles()随后通过runtime.use_fast_path()判断是否启用快速路径并设置 fallback 标志。因此process group 必须在 EP buffer 构造或刷新之前就已初始化且健康。Buffer.dispatch(...)recv_x, recv_count, handle, event, hook buffer.dispatch( x, topk_idx, active_ranks, num_max_dispatch_tokens_per_rank, num_experts, timeout_us, use_fp8True, async_finishFalse, return_recv_hookFalse, )参数x本地 token 隐状态形状[num_tokens, hidden]在 CUDA 上通常为 BF16。topk_idx被选中的专家 ID形状[num_tokens, top_k]用-1标记被掩码的选择。active_ranksEP 级 rank 健康张量形状[num_ranks]dtypetorch.int32。超时检测可能将失败的源 rank 置0。num_max_dispatch_tokens_per_rank每个 rank 的工作区容量应至少等于当前步中跨 rank 的最大本地num_tokens。num_experts全局专家数必须能被num_ranks整除。timeout_us微秒级超时传-1禁用超时检测。use_fp8为True时 dispatch 返回 FP8 数据加 scales。async_finish为True时返回的张量会关联到返回的 event用于流生命周期管理。return_recv_hook为True时调用返回的hook()完成接收同步否则使用event.current_stream_wait()。返回值recv_x打包的本地专家输入。若use_fp8True为(packed_data, packed_scales)元组否则为 BF16 张量。recv_count每个本地专家收到的 token 数。handlecombine()与get_next_combine_buffer()所需的不透明元数据。eventEventOverlap辅助对象不使用 hook 时消费输出前调用event.current_stream_wait()。hook当return_recv_hookTrue时返回的可选同步 hook。从源码看dispatch 输出是本地专家主序local-expert-major的打包布局packed_recv_x形状为[num_local_experts, num_ranks * num_max_dispatch_tokens_per_rank, hidden]同时产出packed_recv_count、packed_recv_src_info与packed_recv_layout_rangehandle正是由(packed_recv_src_info, packed_recv_layout_range, num_max_dispatch_tokens_per_rank, hidden, num_experts)组成的元组。断言还要求x为连续二维 BF16 张量、topk_idx为连续 int64、x.size(1)能被 16 与 128 整除、num_experts % group_size 0。Buffer.combine(...)combined_x, event, hook buffer.combine( x, topk_idx, topk_weights, active_ranks, timeout_us, handle, zero_copyFalse, async_finishFalse, return_recv_hookFalse, outNone, )参数x按dispatch()返回布局打包的本地专家输出。topk_idx与topk_weights将专家输出合并回本地 token 的路由元数据。active_ranks与dispatch()使用的同一个 EP 级 rank 健康张量。timeout_us微秒级超时传-1禁用超时检测。handle匹配的dispatch()调用返回的 handle。zero_copy为True时将专家输出写入buffer.get_next_combine_buffer(handle)再把该张量传给combine()。out合并结果的可选输出张量。Buffer.get_next_combine_buffer(handle)返回用于零拷贝专家输出的下一个 combine buffer。仅可与匹配的 dispatchhandle一起使用并把返回的张量以zero_copyTrue传回combine()。注意同一个 handle 不能跨无关的 dispatch/combine 对复用。Buffer.update_ep_member()在后端成员关系变化后重新连接 EP 对端。PG 恢复更新 rank 活跃度后调用它EP 的传输元数据与 QP 即可刷新。源码中它等价于以is_updateTrue重新执行connect()期间会先调用runtime.update_local_qpns()再重新交换对端元数据。流同步模型EP 操作返回EventOverlap对象并可选择返回 hook见 python/mooncake/mooncake_ep_buffer.py 中的EventOverlap若return_recv_hookFalse在消费输出张量前调用event.current_stream_wait()。若return_recv_hookTrue在选定的重叠点调用返回的hook()。若async_finishTruewrapper 会在 event 辅助对象中额外记录张量保证异步使用与 CUDA graph 场景下的生命周期安全。当把 EP 与自定义专家 kernel、CUDA graph 或应用级流组合时务必显式完成同步。默认 NCCL 后端ElasticBuffermooncake.mooncake_elastic_buffer.ElasticBuffer现在默认使用transportauto。自动模式在扩展以 NCCL Device API 构建且推断出的 EP 拓扑受已编译 NCCL kernel 支持时使用 NCCL现有构造函数调用无需任何修改。若 NCCL 不可用自动模式会回退到 IPC IBGDA 并保留原有后端行为。与之前一样所请求的负载必须具有已编译的 elastic kernel 形状。NCCL 支持是 opt-in使用-DWITH_EPON -DUSE_CUDAON -DUSE_NCCL_DEVICEON构建该选项默认OFF。启用 NCCL 的 EP 扩展目前直接链接libnccl因此即使未选择 NCCL 传输导入mooncake.ep也需要匹配的 NCCL 运行时对于必须兼容旧 NCCL 运行时的部署请保持该选项关闭。无需应用侧 communicator bootstrap。自动模式在 process group rank 0 上生成一个 NCCL unique ID 并通过 group 广播import torch.distributed as dist from mooncake.mooncake_elastic_buffer import ElasticBuffer # Run this program with torchrun so rank metadata is available. dist.init_process_group(backendnccl) buffer ElasticBuffer( dist.group.WORLD, num_max_tokens_per_rank128, hidden4096, num_topk8, ) print(fMooncake EP selected {buffer.transport}) try: # Call buffer.dispatch(...) and buffer.combine(...). pass finally: # Deterministic collective cleanup is recommended when NCCL was selected. buffer.destroy() dist.destroy_process_group()从 python/mooncake/mooncake_elastic_buffer.py 的源码可见_requested_transport()支持auto/ibgda/nccl三种取值auto还会读取环境变量MOONCAKE_EP_TRANSPORT进行滚动发布覆盖_nccl_topology_supported()则按num_ranks num_rdma_ranks * num_nvlink_ranks校验拓扑形状。受控发布时可显式传transportibgda或设置MOONCAKE_EP_TRANSPORTibgda显式transportnccl会禁用自动回退在 NCCL 支持不可用时直接报错。explicitly_destroy参数仍为兼容 DeepEP API 而保留为可选集体调用destroy()仍是销毁 process group 前释放 NCCL symmetric window 最可预期的方式。NCCL 后端的约束需要 NCCL 2.30.4 或更新版本且带 Device API 与 GIN 支持。构建 Mooncake 所用的 NCCL 头文件必须与加载的libnccl完全一致NCCL 升级后需重建 Mooncake。若 PyTorch 会先加载另一个 NCCL需在初始化 process group 前配置或预加载匹配的运行时。Process-group rank 必须构成连续、等规模的 NCCL LSA team。已编译 kernel 支持 1 个 2 或 8 GPU 的 team1x2或1x8、2 个 4 或 8 GPU 的 team2x4或2x8以及 4 个 4 GPU 的 team4x4跨 team 通信使用 hybrid 模式与 rail GIN。其他形状下自动模式会选择 IPC IBGDA。超过一个 rank 的 group 都会请求 GIN 资源包括数据路径仅停留在单个 LSA team 内的运行。逻辑成员保持固定。在每个进程都观察到 Mooncake PG 已恢复所有原始逻辑 rank 槽位后在 EP 迭代之间进行重建每个幸存进程在其既有 buffer 上调用update_ep_member()每个替换进程则用相同参数构造ElasticBuffer。无进程被替换的健康重配置对每个 rank 调用update_ep_member()。这些调用是一次协调操作不要在全部调用返回前开始新的 EP 操作。重配置会重建 host/device communicator、GIN 资源、symmetric window 与 buffer 分配同时保留幸存者的 Python buffer 对象恢复后的 placement 仍须匹配受支持的拓扑。更新前创建的所有 dispatch handle 与 view 均失效不得复用。替换就绪前更新会临时持有两代完整 NCCL 资源communicator、symmetric window、EP buffer 分配、独占 GIN context接近内存、GIN-context 或 QP 上限的部署必须为两代资源预留容量。update_ep_member()不会完成被故障中断或缺少逻辑 rank 的操作PG 恢复后需重试被中断的工作。内部状态 collective 建立之前发生的 rank 本地故障例如配置/运行时不匹配、无法分配最小 CUDA 控制资源无法原地恢复可能需要重启 process group。请在所有 rank 上使用一致的 NCCL/CUDA 配置。Active-rank 张量PG 与 EP 的区分API 表面存在两个 active-rank 张量PG active-rank mask传给pg.MooncakeBackendOptions镜像 Coordinator 已提交的成员关系。EP active-rank tensor传给Buffer.dispatch()与Buffer.combine()。它同样是 rank 级[num_ranks]、torch.int32并且可能被 EP kernel 在超时检测将某个对端标记为失败时更新。在简单集成中它们的取值可能恰好一致但语义不可互换PG 成员关系是配置而 EP 可能从 kernel 级超时观测更新自己的 mask。在把已提交的 PG 成员关系传播到 EP 时要保持 mapping、dtype、device 与容量的一致性。另外注意EP 级超时检测可以把某个 rank 标记为不活跃但更高层仍需协调恢复与路由决策——故障时需同时更新 PG active-rank 状态停止 collective 等待失败 rank、调度/MoE 路由停止向不可用专家分配 token、替换 rank 的 PG 弹性加入以及 EP buffer 的update_ep_member()刷新。测试与示例仓库提供了覆盖 PG 集合、弹性恢复以及 EP 正确性与故障仿真的测试均可作为 API 的可执行契约PG collectivesmooncake-pg/tests/test_pg_collectives.pyPG 弹性恢复与子 group 扩展mooncake-pg/tests/test_pg_elastic.pyPG 基准工具mooncake-pg/benchmark/README.mdEP 正确性与故障仿真python/tests/ep/test_ep_grid.pyEP wrapper 示例python/tests/ep/test_mooncake_ep.pyNCCL EP rank 替换恢复python/tests/ep/test_elastic_buffer_recovery.py在仓库根目录、具备两个可见 CUDA 设备且安装 NCCL 版 EP/PG 扩展的条件下运行 NCCL EP 恢复测试python -m pytest -q python/tests/ep/test_elastic_buffer_recovery.py这些测试覆盖 worker 替换前后的 dispatch/combine、拒绝过期 EP handle以及预留 PG 容量与mooncake-cpu控制组控制组使用 CPU 张量时 EP 数据路径仍保留在 GPU 上。测试复用 PG worker harness并通过MOONCAKE_PGTEST_DEVICE_FILTERS指定 NIC 选择。PG 侧的测试入口还包括 mooncake-pg/tests/test_pg_init_functional.py初始化/单 rank/子 group/销毁/再初始化、mooncake-pg/tests/test_pg_p2p.py直接与批量化 P2P、乱序、多发送方、故障检测、mooncake-pg/tests/test_pg_inference_topologies.pyTP/PP/DP/EP 与 prefill/decode 布局等可用以下命令从仓库根目录运行# 运行全部 PG 测试 python -m unittest discover -s mooncake-pg/tests -v # 仅 CPU 后端 PG 测试 python -m unittest discover -s mooncake-pg/tests -k CPU -v # 仅 CUDA 后端 PG 测试 python -m unittest discover -s mooncake-pg/tests -k CUDA -v # 集合通信基准冒烟 PYTHONPATHmooncake-pg \ python mooncake-pg/benchmark/pgbench.py \ --collective all_reduce --backend mooncake --device cuda -g 2 -b 8 -e 1M -f 2常见问题排查要点导入报 PyTorch 版本错误Mooncake PG was not built against torch...先通过python -c import torch; print(torch.__version__)确认版本再安装匹配该 PyTorch/CUDA 组合的 wheel或以cmake .. -DWITH_EPON make -j重新构建并确保运行时与构建使用同一 Python 环境。dist.get_world_size()与活跃 rank 数不一致这是动态槽位模型的正常表现——max_group_size是固定槽位容量dist.get_world_size()是最高活跃 in-group rank 1空洞可见但被 active mask 跳过用pg.get_active_ranks(backend)检查当前 mask。join_group()挂起或激活超时常见原因是没有任何在线 rank 提交activate_ranks()/recover_ranks()或未来活跃集合并非双向连通导致激活提议挂起至超时。EP 意外把 rank 标记为不活跃先调大timeout_us排除慢启动/调度抖动检查源 rank 是否退出或跳过匹配的 dispatch/combine 调用确认num_experts % num_ranks 0且所有 rank 使用一致的num_experts、top_k、buffer 尺寸假设并保证所有 rank 按相同顺序调用 dispatch 与 combine测试时可传timeout_us-1关闭超时检测。EP 回退而非走快速路径检查容器内 RDMA 设备/驱动可用性、HCA 选择、GPUDirect RDMA / peer-memory 支持、GPU 间 CUDA IPC / peer access以及构建时是否包含所需加速器支持可先用pg.set_device_filter([...])限制 HCA再先跑通 PG collectivesEP 元数据交换依赖健康的 process group。EP 输出不匹配或 buffer 溢出通常是num_max_dispatch_tokens_per_rank小于实际每 rank token 数、num_experts不可被num_ranks整除、topk_idx含超出[0, num_experts)的专家 ID掩码-1除外、handle 被无关 combine 复用或zero_copyTrue未配合get_next_combine_buffer(handle)使用。更多细节可参考 PG/EP 故障排查指南提交 bug 报告时建议附上 Mooncake 提交号与安装方式、PyTorch/CUDA 版本、后端名称mooncake或mooncake-cpu、world size / max world size / rank / 子 group 布局、active-rank 张量的 dtype/device/值、HCA filter 以及各 rank 初始化与故障检测阶段的日志。赞分享人工智能大模型模型推理服务后端【免费下载链接】MooncakeMooncake is the serving platform for Kimi, a leading LLM service provided by Moonshot AI.项目地址https://gitcode.com/gh_mirrors/mo/Mooncake点击查看免费下载相关推荐如何利用Mooncake实现弹性专家并行MoE模型推理的容错与扩展性如何利用Mooncake实现弹性专家并行MoE模型推理的容错与扩展性 MoE专家混合模型凭借其卓越的扩展性和推理效率正成为大语言模型发展的关键方向。然而人工智能大模型模型推理服务后端torchtitan-npu EP 通信域分离为 Ascend NPU 上的 MoE 专家并行创建独立 HCCL 通信组torchtitan npu EP 通信域分离为 Ascend NPU 上的 MoE 专家并行创建独立 HCCL 通信组 本指南讲解 torchtitan n人工智能大模型分布式训练预训练模型优化Ascend25MB 数据库管理工具 DBX一个客户端连 100 数据库25MB 数据库管理工具 DBX一个客户端连 100 数据库 你电脑上是不是也躺着好几个数据库客户端MySQL 用一个、Redis 用一个、MongoDB数据库客户端数据库桌面应用CLI后端MCP 服务AI 应用上一篇TV-Bro 免费开源电视浏览器让遥控器真的能操作网页下一篇在 Windows 上装安卓应用只要 3 步APK Installer 完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表