
EnergonAI的核心AsyncEngine如何将多GPU分布式推理封装成单次调用【免费下载链接】EnergonAILarge-scale model inference.项目地址: https://gitcode.com/gh_mirrors/en/EnergonAIEnergonAI 是一个面向大规模模型推理的开源服务框架它最核心的设计就是把多GPU分布式推理张量并行 流水线并行封装成「提交任务 → 拿到结果」的单次调用。你不需要理解进程怎么拉起、GPU 之间怎么通信只需调用两行代码就能让一个 175B 参数的大模型跑在多台 GPU 上 。本文带你快速看懂AsyncEngine背后的完整机制。为什么需要封装多GPU推理推理一个超大模型往往单卡放不下必须拆分张量并行TP把每一层的权重横向切开分布在多张 GPU 上前向时多卡同时计算再聚合流水线并行PP把模型按层切开不同 GPU 负责不同的阶段像流水线一样接力处理。如果让你自己写要处理进程管理、RPC 注册、设备映射、阶段间数据搬运……代码量庞大且容易出错。EnergonAI 用AsyncEngine解决了这个问题它把「单实例多设备」SIMD的分布式执行伪装成「单实例单设备」SISD的普通函数调用用户视角里模型就是一个对象。架构全景一个主控 N 个GPU工作进程整个引擎由两部分组成启动入口是一行 launch_engine┌────────────── master 进程主程序 ──────────────┐ │ submit(uid, data) │ │ │ 进入 submit_queue │ │ ▼ │ │ ┌─────────────┐ ┌──────────────────┐ │ │ │ _submit_loop│───▶│ _completion_loop│◀───┐ │ │ │ (打包批处理) │ │ (收集结果按uid拆分) │ │ │ │ └──────┬──────┘ └────────┬─────────┘ │ │ └─────────┼────────────────────┼────────────┼──────┘ │ send (RPC Pipe) │ recv │ ▼ ▼ │ ┌───────────┐ ┌───────────┐ ┌───────────┐ │ │ worker 0 │───▶│ worker 1 │───▶│ worker N-1│──────────┘ │ GPU: TP组/ │ │ GPU: TP组/ │ │ GPU: 最后 │ │ PP第0阶段 │ │ PP第1阶段 │ │ 一个PP阶段 │ └───────────┘ └───────────┘ └───────────┘各组件分工组件位置职责AsyncEngineenergonai/engine.py主控提交队列、两条后台线程、结果映射Workerenergonai/worker.py每 GPU 一个进程加载模型分片、执行前向Pipeenergonai/pipe.py基于 RPC 的命名队列master 与 worker 间传数据BatchManagerenergonai/batch_mgr.py请求打包策略默认逐条可扩展动态批处理launch_engineenergonai/engine.py一行代码拉起全部 worker 并返回引擎启动时launch_engine会按tp_world_size × pp_world_size用多进程spawn逐个拉起Worker每个 worker 在 worker.py 中注册 ColossalAI 的 1D 张量并行组与流水线组并挂上 TensorPipe RPC默认关闭 SHM 传输以避免超时问题。随后各 worker 通过Pipe与 master 握手master 只认流水线第 0 阶段的 worker 作为任务入口、最后阶段的 worker 作为结果出口。一行调用submit 提交wait 取结果对用户来说分布式推理只剩三步真实用法见 opt_fastapi.pyengine launch_engine(tp4, pp1, ..., model_fnopt_175B, checkpointckpt.pt) engine.submit(uid, inputs) # ① 提交立即返回不阻塞 output await engine.wait(uid) # ② 异步等结果或 engine.get(uid) 同步等submit(uid, data)把请求放进submit_queue立刻返回不阻塞你的服务线程await wait(uid)协程版轮询每 0.1s 查一次结果表天然适配 FastAPI 等异步框架get(uid)则是同步版轮询如果队列设了上限queue_size队列满时会抛出QueueFullError方便上层做 406 限流——这是一种优雅的反压机制 。引擎内部两条后台线程干完所有脏活AsyncEngine构造时会启动两个常驻线程见 engine.py1. 提交线程_submit_loop—— 负责打包发货持续检查submit_queue有数据就交给BatchManager.make_batch()打包成TaskEntry再通过submit_pipes广播到流水线入口。默认的 BatchManager 逐条处理生成类任务可换成动态批管理器如examples/opt/batch.py中做 left padding 桶式分组的策略把多个短请求凑成一批显著提升 GPU 利用率。2. 完成线程_completion_loop—— 负责收货拆单从各流水线的最后一个阶段轮询收集输出所有流水线副本都到齐才算一批完成调用split_batch()把批量输出按uid拆开写进completion_map。提交线程和完成线程都记录耗时日志中会打印batch size与耗时方便你观察吞吐。Worker 侧每个GPU就是一个流水线工位每个Worker的逻辑极简就是一个死循环见 worker.py从input_pipe非阻塞取任务第 0 阶段来自 master其余阶段来自上一段 worker在torch.inference_mode()下执行前向——张量并行的多卡集合通信由 ColossalAI 在模型内部透明完成把结果送进output_pipe最后一段则送回 master。关闭同样是一键的engine.shutdown()会向所有 worker 广播终止信号Terminator并优雅 join 掉后台线程连 CtrlC 都被注册成了自动关机 。快速上手5 分钟跑起 OPT 多GPU推理服务克隆仓库https://gitcode.com/gh_mirrors/en/EnergonAI执行pip install -r requirements.txt pip install .打开示例目录examples/opt/在opt_config.py中设置model_class、checkpoint 路径并把tp_init_size设为你的 GPU 数量运行bash server.sh浏览器打开http://[ip]:[port]/docs即可在线体验推理接口。BLOOM、GPT、BERT、ViT 等更多示例都在examples/目录下结构几乎一致定义模型 调用launch_engine 套一层 FastAPI。总结为什么说 AsyncEngine 是关键抽象屏蔽分布式细节进程拉起、RPC 注册、张量/流水线并行、设备映射全部封装在energonai/内部异步非阻塞submit/wait模型让推理服务能并发处理大量请求批处理收益自动叠加可扩展批策略BatchManager、并行度tp/pp、背压queue_size都是构造参数按需调整即可。理解了AsyncEngine你就掌握了 EnergonAI 的核心它把「多GPU分布式推理」这件事压缩成了每个开发者都能轻松上手的单次调用⚡。【免费下载链接】EnergonAILarge-scale model inference.项目地址: https://gitcode.com/gh_mirrors/en/EnergonAI创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考