目录
文章目录
- 目录
- vLLM 核心技术
- PageAttention
- 连续批处理
- 软件架构
- LLMEngine 和 AsyncLLMEngine
- EngineCore
- Scheduler
- 调度单元 SequenceGroup
- 调度策略
- KVCacheManager
- BlockPool
- Executor
- Worker
- CacheEngine
- ModelRunner
- 核心业务流程
- 加载模型与预分配显存流程
- 加载模型
- 预分配显存
- 调度与推理流程
- 调度
- 推理
- v1 版本的 Offline batching 完整推理流程
- v1 版本的 Online serving 完整推理流程
- Executor-Workers 处理流程
vLLM 核心技术
vLLM 是 UC Berkeley 等开源的一款高吞吐的、内存高效的大模型推理与服务引擎。它的核心创新在于 PagedAttention 算法,将注意力层的 KV cache 以类似操作系统虚拟内存分页的方式管理。通过这一创新,vLLM 实现了远超传统方案的推理性能。
如下图。一个 13B 的模型在 A100 40GB 上做推理时的显存分配情况,可以看见 KV cache 占用较高。因此,如何优化 KV cache,节省显存,提高吞吐,就成了 LLM 推理引擎需要解决的重点问题。vLLM 在显存使用率和推理吞吐方面都有明显的提升。
PageAttention
在 LLM 推理时,为了不重复计算历史信息,系统会缓存每个 token 的 Key 和 Value 向量,这就是 KV Cache。随着对话进行,KV Cache 会动态增长。如果存储这些重要的 KV cacahe 是所有推理引擎需要解决的首要问题。
在 PageAttention 之前的传统推理框架中,当收到一条请求时,它会为这条请求中的 prompts 分配 GPU 显存空间,包括对 KV cache 空间的分配等。但由于请求 prompts 的 seq_len 是无法预定的,所以会按照 [batch_size, max_seq_len] 中 max_seq_len 这样的固定最大尺寸来分配。
显然的,传统方式这会造成浪费,这些浪费严重限制了同时处理的请求数量(Batch Size),从而拉低了 GPU 的整体吞吐量。
- 内部碎片:为应对可能的最大长度,系统会过度预留空间,但实际生成长度往往短得多,导致预留空间大量闲置。实验表明,在传统系统中,KV Cache 的实际有效利用率可能低至 20.4%。
- 外部碎片:不同请求释放后留下的内存空洞,因大小不一而难以被后续请求利用。
Paged Attention(页面注意力)是 vLLM 推理引擎的核心技术(https://arxiv.org/pdf/2309.06180)。PageAttention 借鉴了 Linux 的 “虚拟内存分页技术” 为 GPU 显存也实现了 “虚拟显存分页技术”,以此来管理 KV Cache。
- 分页存储:定义 “虚拟显存块” 和 “物理显存块” 的概念,并使用 Block table 来建立映射。将每个序列的 KV Cache 分割成固定大小的 Block,每个 Block 包含固定数量(如 16 个)token 的 KV 向量。
- 非连续存储:KV Cache 被存储在这些 Block 中。逻辑上,这些 Block 是连续的;但物理上,这些 Bloack 是非连续的,可以分散在物理显存的任意位置,这样就避免了物理显存的碎片。
- 按需分配:当序列产生新 token 的 KV cache 时,先尝试放入该序列的最后一个 Block,如果该 Block 满了,则申请一个新的物理块并映射为序列的下一个逻辑块。整个推理过程按需分配,不进行预留,减少显存浪费。
- 高效释放:当请求结束或被中止,其占用的物理块会立即释放,且无需大范围的释放操作,减少了释放和碎片整理的开销。由于所有块大小相同,释放后很容易被其它序列重新利用,也不会产生严重碎片。
处理单个请求时候的 Block table 映射示意图如下。
处理多个不同请求时,Block table 映射示意图如下。
处理多个相同请求时,Block table 映射示意图如下。可见,PageAttention 同样实现了 “–enable-prefix-caching 公共缓存(共享物理块)” 和 “COW(写时复制)”。
连续批处理
推理引擎是按 batch(批次)来处理请求的,因为这样满足 Input Tensor [batch_size, seq_len, hidden_size] 的并行处理特性。简单来说一个 batch 会包含多个 seqs,并且这些 seqs 可能来自若干个 reqs。
传统推理引擎采用 “静态批处理”,即:收集多个 seqs => 组成一个 batch => 一次性输入 LLM => 所有 seqs 同时计算 => 所有 seqs 完成后,释放整个 batch,再处理下一批。显然的,静态批处理的显存利用效率和吞吐都很差。如下图,seqs 长短不一的请求,短的等长的,就造成了显存浪费和吞吐下降。
vLLM 实现的 “连续批处理” 是一种动态调度技术,在每个推理阶段的粒度(一个前向计算)上进行动态的批次调度管理。每次新的迭代,Scheduler 都会及时将已经完成的 seq 释放点,然后整理当前所有未完成的 seqs 组成新的 batch。通过这种方式,vLLM 的推理流水线能够持续地接纳新 seqs,大幅提升了显存利用率和整体吞吐量。
软件架构
LLMEngine 和 AsyncLLMEngine
vLLM 提供了 2 种 User Interface(客户端调用方式):
- Offline Batched Inference(离线批处理服务),同步调用。
- API Server For Online Serving(在线推理服务),异构调用,支持 2 种 API:
- Simple Demo API Server:测试开发用。
- OpenAI-Compatible API Server:兼容了 OpenAI 请求格式,包括:OpenAI Completions API 和 OpenAI Chat API。
这两种 User Interface 的底层由同一个 LLMEngine 推理引擎核心模块支持,如下图所示。
LLMEngine 提供了以下关键同步接口:
- add_request():将每一个请求包装成 vLLM 能处理的 SequenceGroup 数据类型,并将其加入 Scheduler 的 waiting queue 中。
- abort_request():用于主动终止一个请求。
- step():执行 1 个推理阶段(1 次 Prefill + 2 次 Decode,这样算 3 个推理阶段)。Scheduler 会决定要送那些数据去执行本次推理,并负责给这些数据分配好物理块。
AsyncLLMEngine 是对 LLMEngine 的异构调用包装,提供了以下关键异步接口:
- generate():将同步的 add_request() 重写成异步的形式。
- abort():见同步的 abort_request() 重写成异步的形式。
EngineCore
Scheduler
Scheduler 的主要作用是在 1 个推理阶段中,完成 seqs 的调度和抢占(决定要把哪些数据送给模型做推理),同时负责分配好 KV Cache 的物理块 id,并做好逻辑块和物理块的映射(逻辑层面的显存分配)。注意,这里只是分配了物理块的 id,而不是物理块本身。因为物理块的实际分配是模型在推理阶段中由 CacheEngine 根据物理块 id 来完成的。
- self.policy:根据调度策略(FCFS),对各个队列里的 seq_group 按照其 arrival time 进行排序。
- self.prev_time:上一次调度发起的时间点,初始为 0。
- self.prev_prompt:取值为 True or False,初始为 False。若上一次调度时,Scheduler 有从 waiting 队列中取出 seq_group 做推理,即为 True,否则为 False。
- self.last_prompt_latency:记录 “当前调度时刻 - 最后一次有从 waiting 队列中取数做推理的那个调度时刻” 的差值,初始为 0。并不是每一次调度都会从 waiting 队列中取 seq_group,它可能依旧继续对 running 队列中的数据做推理。
- BlockManager:物理块管理器,维护着两个重要属性:
- BlockAllocator:物理块分配者,负责实际为 seq 做物理块的分配、释放、拷贝等操作。其下又分成 self.gpu_allocator 和 self.cpu_allocator 两种类型,分别管理 GPU 和 CPU 上的物理块。
- self.block_tables:负责维护每个 seq 下的 Block table。
调度单元 SequenceGroup
在实际的推理场景中,一个 request、一条 prompt 可能对应多个 seqs。例如下列 2 种 Decode 策略:
- Parallel Sampling(并行采样):输入 1 条 prompt,模型返回 n 种不同的 seqs。
- Beam Search(束搜索):输入 1 条 prompt,每个推理阶段都会返回 topK 个 seqs,其中 K 被称为 Beam width(束宽)。另外,如下图可以看见每个推理阶段(生成一个新 token)把 topK 个序列喂给模型时,它们的前置 token 中有大量的 KV cache 是重复的。这里又利用了 PageAttention 共享物理块的优势。
因此,为了管理 request、prompt、seqs 之间的复杂关系,vLLM 定义了 SequenceGroup 来进行封装。一个 SequenceGroup 包括了 1 个 request、1 条 prompt 以及多个 seqs。
- self.seqs_dict:每个 seq 是一个 Sequence 对象,一个 seq_group 下包含若干个 seqs。
- self.sampling_params:采样参数。
- get_max_num_running_steps:该 seq_group 在剩余生命周期内并行 running 的最大 seqs 数量。
- Sequence 对象的 self.logical_token_blocks 属性:每个 seq 都单独维护一份属于自己的逻辑块,不同的逻辑块可以指向同一个物理块。
- Sequence 对象的 _append_tokens_to_blocks 方法:没有逻辑块就分配,有逻辑块就看最后一个逻辑块的插槽是否已满,没满就插入,满了就增加一个新的逻辑块。
每个 SequenceGroup 对象都拥有以下状态类型:
- WAITING:处于 waiting 队列,该队列中的 seqs 都没有做过 Prefill。
- RUNNING:处于 running 队列,该队列中的 seqs 都送去了 1 次推理阶段。
- SWAPPED:处于 swapped 队列,此时 GPU 显存不足,处于后方的 SequenceGroup 被抢占,推理暂停,相关的 KV block 被 Swap 到 CPU 主存中了(Swap out),等待 GPU 显存资源充足后再置换回来重新计算(Swap in)。
- FINISHED_STOPPED:正常执行完毕,例如 EOS 符号了。
- FINISHED_LENGTH_CAPPED:因生成的 output_seqs 触达长度上限而结束推理。
- FINISHED_IGNORED:因输入的 prompt 过长,请求被直接忽略了。
- FINISHED_ABORTED:因不正常状态而被终止推理,例如客户端断开连接了。
调度策略
running 队列中的 seq_group 不一定能继续在本次调度中被选中,这是因为 GPU 显存使用情况一直在变,以及 waiting 队列持续有新的请求进来。所以调度策略的职责就是要根据这些变动,对送入模型做推理的数据做动态规划。总结来说:
如果当前 swapped 队列为空,那就去检查是否能从 waiting 队列中调度 seq_group,直到不满足调度条件为止(例如:GPU 空间不足,或 waiting 队列已为空等)。此时,1 个推理阶段中,所有的 seq_group都处在 Prefill 阶段。
如果当前 swapped 队列非空,或者无法从 waiting 队列中调度任何 seq_group 时。此时,1 个推理阶段中,所有的 seq_group 要么全来自 running 队列,要么来自 running + swapped 队列,它们都处在 decode 阶段。
- 检查是否能从 running 队列中调度 seq_group,直到不满足调度条件为止。
- 若本次无新的被抢占的 seq_group,且 swapped 队列非空,就检查是否能从 swapped 队列中调度 seq_group,直到不满足调度条件为止。
以 swapped 是否非空作为判断入口的原因是 swapped 标识了 GPU 显存资源是否足够。
V0 版本,在 1 个推理阶段中,所有的 seq_group 要么全部处在 Prefill 阶段。要么全部处在 Decode 阶段。而 V1 版本开始允许单次推理阶段中同时调度 Prefill 和 Decode 的请求,同时调度器中只维护 waiting 和 running 队列。
KVCacheManager
KVCacheManager,最早称为 BlockSpaceManager,负责物理块 id 的分配和虚实映射,并且支持 GPU 显存和 CPU 主存的分配。注:之所以需要 CPU 主存,是因为当 GPU 显存不足时,就会把后来的请求先抢占掉,把其相关的 KV Cache 先 Swap(置换)到 CPU 主存上,等后续 GPU 显存充足了,再把它们加载回来。所以同样需要对 CPU 主存进行管理。
KVCacheManager 主要负责管理 “req 和它们的 Blocks”:
- BlockPool:这张 GPU 卡上所有的物理块;
- req_to_blocks:一个 req 所拥有的所有的物理块;
- req_to_block_hashes:一个 req 所拥有的每一个物理块的 hash 值。
- num_cached_block:一个 req 下所有满块的数量。
上图中每个绿色块就是一个 KVCacheBlock 实例。假设一块 GPU 上的物理块共有 N 个,那么这里就有 N 个 KVCacheBlock 实例,但是这些实例并没有真正存储 KV Cache 数值,这是分配了 Block id。可见,KVCacheManager 只负责区域划分和分配,并不真正执行推理。
BlockPool
KVCacheBlock、FreeKVCacheBlockQueue、cached_block_hash_to_block 共同组成了 BlockPool 实例,也就是一张 GPU 上的 Block 池,它负责维护这张卡上所有的物理块。
Executor
Executor 是 Worker 的控制面,Worker 是 GPU 设备的抽象,一个 Worker 就是一块 GPU。
Executor 可以指定使用什么方法来控制这些 Workers、也负责分布式环境的初始化(支持 TP、PP 并行推理),目前支持的方法有如下 4 种,通过 --distributed-executor-backend 指定。默认为 None,即:根据实际的分布式配置(world_size)和平台特征(cuda、是否使用 Ray 分布式计算框架等)来自动选择。
- mp(MultiprocExecutor):适用于单机多卡场景。当单机的卡数满足分布式配置,且没有正在运行的 Ray pg 时,默认使用。此时 Executor 是一个主进程,其下有若干个 workers 子进程。
- ray(RayDistributedExecutor):适用于多机多卡。安装了 Ray 且分布式配置需要多机时使用。此时 Executor 成为一个 Ray driver process,管理着若干个 worker process。(–distributed-executor-backend ray)
- uni(UniProcExecutor):适用于单卡场景。
- external_launcher(ExecutorWithExternalLauncher):想要用自定义的外部工具(如 Slurm)来做分布式管理。
Worker
Worker 的作用是将模型参数 Load 到 GPU 上,然后对 Scheduler 传来的数据执行 1 次推理,最后返回结果。
CacheEngine
CacheEngine 负责管理 GPU 显存和 CPU 主存上的 KV cache 物理。注意,Scheduler 下的 BlockManager 只负责分配物理块 id;而 CacheEngine 则根据 id 分配实际的物理块。
ModelRunner
ModelRunner 负责加载模型,并执行推理。PagedAttention 的相关逻辑在这里实现。
核心业务流程
加载模型与预分配显存流程
启动 vLLM 程序进行推理之前,首先要完成模型的加载和 GPU 显存的预分配。
加载模型
执行以下代码时候,vLLM 把 model 加载到 workers 上。
llm=LLM(model="facebook/opt-125m")预分配显存
实例化了一个 LLMEngine 对象时候会对 GPU 显存和 CPU 主存进行预分配。LLMEngine 通过模拟实验(profiling)的方式来决定 GPU 显存和 CPU 主存上到底有多少个 KV cache 物理块可用于分配的。这个步骤称为profile_num_available_blocks,包括:
- 杜撰假数据:两个重要参数。
- max_num_batched_tokens:指示了 1 个推理阶段中,LLMEngine 能处理的最大 token 数量。默认是 2048。
- max_num_seqs:指示了 1 个推理阶段中,LLMEngine 能处理的最大 seq 数量。默认是 256。(注意,1 条 seq 指待推理的 1 条数据。)
根据这两个参数,假设在推理中,平均一个 seq 要处理 max_num_batched_tokens // max_num_seqs 个 token,余数部分我们默认放在第一个 seq 中。那么当 max_num_batched_tokens=10、max_num_seqs=3 时,就需要杜撰出 3 条 seq 假数据,它们的长度分别为 4、3、3。
- 用假数据模拟一次前向推理:在 1 个推理阶段中需要分配多少的显存给 KV cache 可以使用如下公式计算。其中,可以用杜撰出来的假数据模拟一次前向推理得到 “不使用 KV cache 做 1 次推理时的显存占用”,继而求得 “分配给 KV cache 的显存”。
分配给 KV cache 的显存=GPU 总显存 - 不使用 KV cache 做1次推理时的显存占用(包括模型本身和推理阶段中的中间数据)- 计算可分配的 KV cache 物理块总数:知道分配给 KV cache 的显存总量后,就可以计算总的物理块数量了。公式如下,其中重要参数有:
- slot_size:指示了 Slot 的空间大小。从下述公式可知,slot_size 就是一个 token 词向量的大小。
- block_size:指示了 一个物理块有多少个 Slot。
slot_size=num_heads*head_size*dtype K_cache_block_size=V_cache_block_size=block_size*slot_size 总物理块数量 num_blocks=分配给 KV Cache 的显存大小 / K_cache_block_size(单位 Bytes)# kv_cache_shape: (2, num_blocks, block_size * num_kv_heads * head_size)注意,CPU 主存不需要做模拟实验,因为分配给 vLLM 的 CPU 主存是用户主动传参控制的,默认是 4G。也意味着只能在这 4G 上做 Swap。同理,将上面公式中 “分配给 KV Cache 的显存大小” 替换成 4G,就能得到 CPU 上物理块的数量。
- 将预分配的 KV Cache 加载到 GPU 显存上:确定好 KV Cache block size 后,就开始创建 Empty tensor 用于放置到 GPU 上完成显存的预分配。以后这块显存就是专门用来做 KV Cache 的了。
调度与推理流程
以 Offline Batched Inference 为例,当调用 llm.generate(prompts, sampling_params) 时,实际做了 2 件事情:
- add_request() 调度:将输入的数据(prompts 和 sampling_params)传给 LLMEngine,把每 1 条 prompt 都封装为 1 个 SequenceGroup 对象,然后把 SequenceGroup 加入到 Scheduler 的 waiting 队列等待处理。(注:vLLM 将一个 Prompt 视为一个请求,对应一个 SequenceGroup。)
- run_engine() 推理:只要 Scheduler 的 waiting、running、swapped 队列非空,那么就会调用 step() 来执行 1 个推理阶段。
fromvllmimportLLM,SamplingParamsif__name=="__main__":# 一个批次prompts=["Hello, my name is","The president of the United States is","The capital of France is","The future of AI is",]# Create a sampling params object.sampling_params=SamplingParams(temperature=0.8,top_p=0.95)# 创建 llm 实例,在这个过程中也创建了 llm 实例下的 llm_enginellm=LLM(model="facebook/opt-125m")# 执行offline batching推理,得到这批 prompts 的输出outputs=llm.generate(prompts,sampling_params)# 打印输出foroutputinoutputs:prompt=output.prompt generated_text=output.outputs[0].textprint(f"Prompt:{prompt!r}, Generated text:{generated_text!r}")调度
vLLM 调度和抢占总原则是 FCFS(先来先服务),后来先抢占,GPU 不够就先 Swap 到 CPU 上。
下面介绍 Scheduler 的调度和抢占流程,假设这个 seq_group 只有 1 条 seq。
- 在推理开始之前,LLMEngine 对 seq 进行分词,封装为 seq_group,然后进行 waiting 队列,seq_group 状态为 waiting;
- 第 1 个推理阶段,Scheduler调度了这个 seq_group,由于它的采样参数中 n=4,所以在做完 Prefill 后,它会生成 4 个 seq,它们的状态都是 running;
- 若干个推理阶段后,GPU 显存资源不够了,这个 seq_group 被抢占了(preemption),此时对它的处理有两种方式:
- seq 数量 > 1:采取 Swap 策略,seq_group 所有的 seqs 及其相关的 KV block 被 swap out 到 CPU 主存上。此时所有 seq 的状态变为 swapped。因为 seq 数量比较多,直接把 KV block 抛弃,比较可惜。
- seq 数量 = 1:采取 Recomputation 策略,把该 seq_group 相关的物理块都释放掉,然后将它重新放回 waiting 队列中。因为 seq 数量少,重新计算 KV block 的成本不高。
- 又过了若干个推理阶段,GPU 线程资源又充足了,此时执行 swap in 操作,继续对 seq_group 做推理,此时 seqs 的状态又变为 running。
- 又过了若干个推理阶段,该 seq_group 中有 1 个 seq 已经推理完成了,状态就被标记为 finish,此后这条已经完成的 seq 将不参与调度。
- 又过了若干个推理阶段,该 seq_group 下所有的 seqs 都已经完成推理了,这样就可以把它作为最终 output 返回了。
推理
Executor 接受到来自 Scheduler 的 seq_group 之后就将 seqs 分发到各个 workers 上去做推理。需要两个模型来完成:
- CacheEngine 负责管理实际的 KV Cache 数据。
- ModelRunner 负责加载模型,进行推理。
v1 版本的 Offline batching 完整推理流程
在这里梳理 Offline batching 的整个运作流程如上图所示。
- process0 是 LLMEngine 进程,核心是 SyncMPClient(Sync MultiProcessing Client),对 req 的 processor 输入、output_processor 输出进行处理。
- process1 是 EngineCore 进程,核心是 Scheduler 和 Executor,对 req 的实际推理过程进行处理。
- process0 和 process1 之间使用 ZeroMQ 进行通信。ZMQ 是一个 高速、异步的消息队列通信库,主要用于进程间(IPC)、服务器之间、线程之间的高性能通信。
可见,vLLM 将 CPU(SyncMPClient) 和 GPU(Scheduler、Executor)两个部分的业务逻辑进行了 ZMQ 的解耦,这样可以更好实现 CPU 和 GPU 之间的 Overlap。
上图所的具体流程为:
- add req to:LLM 接收 req,并将 req 发送至 LLMEngine processor。
- process req:processor 对 req 进行预处理,具体包括:Tokenize、验证输入参数的合法性、将原始 req 封装为 EngineCoreRequest 等等。
- encode & send req:将 EngineCoreRequest 发送给 SyncMPClient 后,SyncMPClient 对其做 ZMQ encode(编码)操作转换为 pickle EngineCoreRequest 格式。然后通过 zmq input_socket 将 pickle EngineCoreRequest 发送给 EngineCore。
- recv & decode req:EngineCore 启动了一个线程,在该线程中执行 process_input_socket() 持续监听 input_socket 是否有 pickle EngineCoreRequest。然后对 pickle EngineCoreRequest 进行 ZMQ decode(解码)得到 EngineCoreRequest,最后将 EngineCoreRequest 放入到 input_queue 队列中。
- add req to scheduler:将 EngineCoreRequest 添加到 Scheduler 的 waiting 队列中。
- one step inference:Scheduler 执行 step() 将 EngineCoreRequest 从 waiting 队列移动到 running 队列中开始 1 个推理阶段。最后返回推理结果 EngineCoreOutputs。
- put output to output_queue:将 EngineCoreOutputs 推理结果放入 output_queue 队列中,准备输出。(注:此处会重复执行 5~7 步,直到 input_queue 和 scheduler 中再无 EngineCoreRequest 对象。)
- encode & send req:EngineCore 还会启动一个线程,在该线程中执行 process_output_socket() 持续监听 output_queue 中是否有 EngineCoreOutputs。然后将 EngineCoreOutputs decode 为 pickle EngineCoreOutputs,最后装入到 output_queue 中。
- process_output:对 EngineCoreOutputs 做一些后置处理,例如 detokenize 等。
- output until all reqs are finished:LLM 回去检测每一条 req 是否已做完推理(req.finished的状态),等到全部的 req 都完成推理后,LLM 会一起输出推理结果。
v1 版本的 Online serving 完整推理流程
和 Offline bacthing 整体流程相比,Online serving 的 EngineCore 流程没有变,所以只需要关注 APIServer 流程的区别。
- async per req generate:API server 接收到 req 并发起 generate(),将 req 发送给 AsyncLLMEngine。
- create asyncio.queue() for each req 和 register the queue to output_processor:都在 AsyncLLMEngine 上完成。
- create asyncio.queue() for each req:针对每一个 req 创建一个单独异步队列,这个异步队列将用于存储这个 req 的输出结果。因为 online serving 的 req/resp 是流式处理的,相较于 offline batching 对所有 reqs 进行一次性统一输出,online serving 就有必要将各个 reqs 的输出隔离开来,所以需要为每个 req 都创建一个异步队列。
- register the queue to output_processor:把这个 req 的异步队列注册到 output_processor 中。
- encode & send req、recv & decode req 等等:流程与 offline batching 基本一致,这里不再赘述。
- recv & decode req:AsyncLLMEngine 上启动了一个异步任务 asyncio.create_task(_run_output_handler()),这个函数的主要作用是异步接受、处理来自 EngineCore 的输出结果。包括:get_output_async() 函数持续监听 output_socket 传递来的 pickle EngineCoreOutputs,过程与 offline batching 也一样。
- split output into slices:将 output_queue 中的数据切成 slices,方便后续做流式输出。然后将 slices 数据放入这个 req 所注册在 output_processor 中的、只属于这个 req 的异步队列中。最后 output_processor 将会以 slices 为维度,对输出数据做诸如 detokenize 等的操作。
- continously yield current output:对于一个 req,AsyncLLMEngine 将会持续从 output_processor 中获取到它当前的输出结果,然后流式的返回给用户,直到这条请求推理完毕。
Executor-Workers 处理流程
Scheduler 和 Executor-Workers 都处于 EngineCore 进程中。一个 Executor 主进程管理着若干 workers 子进程,通常一个 workers 对应一张 GPU。
Executor 负责把 req 广播(broadcast)到各个 workers 上执行。各个 workers 接收到 req 完成实际的推理过程后将推理结果返回给 Executor。
下图是一个 MultiprocExecutor 类型,只展示了 1 个 worker,没有画出全部 workers。
- Scheduler 选出送去推理的 req 并交给 MultiprocExecutor 执行推理。
- 在 MultiprocExecutor 上创建一个 rpc_broadcast_mq 队列,用于存储 Executor 要 broadcast 的 “小数据(<=10MB)”,而 “大数据(>10MB)” 则不会进入此队列,而是通过 zmq socket 进行传输。每条数据都是 (method, data) 的形式,指明了 worker 处理的 data 和处理方法。
- 在 MultiProcExecutor 上还通过 make_worker_process 创建了子进程,每个子进程对应一个 worker 实例,用于从 worker 相应的 worker_response_mq 队列中读取返回数据。
- 每个 worker 实例由几部分组成,包括:
- WorkerWrapper:负责管控一个 worker 的生命周期、所占资源、扩展功能(例如在 RLHF 场景下的一些功能)等等。
- rpc_broadcast_mq 队列:连接到 Executor rpc_broadcast_mq 队列,从中读取(method, data)数据。
- worker_response_mq 队列:用于存放这个 worker 上推理输出结果的队列,连接到 Executor worker_response_mq 队列,Executor 就可以从中获得每个 worker 的输出结果了。
- ModelRunner:一个 worker 实例下维护着一个 ModelRunner 实例,该实例维护着大模型的权重分片(model weights sharding)、这块 GPU 上的 kv_caches、attn_backend 等一系列信息,它负责模型权重的加载,并执行实际的推理过程。
- Executor 和 workers 之间的 ZMQ sockets:用于 Executor 和 workers 之间的进程间通信。有多种 socket 类型:
- ready_socket:worker 向 Executor 发送 ready 信号,以及 worker_broadcast_mq_handler。
- local_socket:直接使用 local_socket 做大数据的输入输出通信。
- 等等
- worker_busy_loop():在 worker 上启动一个 busy loop,持续监听 Executor 发送的数据、做推理、并将推理结果持续返回给 Executor。如此的,这个 worker 就无限运转起来了,除非收到显式终止这个 worker 的信号。