当前笔记顺序Engine(当前llm_engine.py和scheduler.py)-LayersEngine的运行架构llm_engine.py作用解析作为顶层编排者负责“怎么执行”。和调度核心Scheduler.py组合使用。import atexit #Python 内置模块用于注册程序退出时自动执行的函数比如清理资源、关闭进程 from dataclasses import fields #Python 内置的dataclasses模块工具用于获取数据dataclass的所有字段名 from time import perf_counter #Python 内置的高精度计时器比time.time更精准用于统计代码执行耗时 from tqdm.auto import tqdm #进度条库auto自动适配环境比如 Jupyter / 终端用于可视化任务进度 from transformers import AutoTokenizer #Transformers 库的分词器工具用于文本↔token_id 的转换 import torch.multiprocessing as mp #PyTorch封装的多进程模块兼容 GPU 张量共享用于创建多进程执行模型推理 from nanovllm.config import Config from nanovllm.sampling_params import SamplingParams from nanovllm.engine.sequence import Sequence from nanovllm.engine.scheduler import Scheduler from nanovllm.engine.model_runner import ModelRunner class LLMEngine: #主要是注册进程尤其是多显卡情况下注册子进程 def __init__(self, model, **kwargs): #提取Config类的所有字段名用于过滤入参只保留 Config 需要的参数,fields(Config) 会返回 Config 类定义的所有属性如 model, max_num_seqs 等。 config_fields {field.name for field in fields(Config)} #获取所有配置了的键值对类型数据 config_kwargs {k: v for k, v in kwargs.items() if k in config_fields} config Config(model, **config_kwargs) #此处配置Sequence的block可以装多少个token Sequence.block_size config.kvcache_block_size # 存储所有子进程 self.ps [] # 存储同步事件用于协调主进程和子进程 self.events [] #使用 spawn 方式创建进程适合 CUDA ctx mp.get_context(spawn) # 如果张量并行度 (tensor_parallel_size) 大于1就启动相应数量的子进程 for i in range(1, config.tensor_parallel_size): #进程间同步的事件锁等待子进程就绪 event ctx.Event() #创建多个 ModelRunner 子进程用于张量并行 process ctx.Process(targetModelRunner, args(config, i, event)) process.start() self.ps.append(process) self.events.append(event) #主进程Rank 0自己实例化一个 ModelRunner 负责统筹 self.model_runner ModelRunner(config, 0, self.events) #加载预训练分词器 self.tokenizer AutoTokenizer.from_pretrained(config.model, use_fastTrue) #获取结束符 config.eos self.tokenizer.eos_token_id self.scheduler Scheduler(config) #注册self.exit函数当程序终止时自动执行比如关闭子进程、释放模型资源 atexit.register(self.exit) def exit(self): #通知主进程的 ModelRunner 优雅退出 self.model_runner.call(exit) #释放主进程的 ModelRunner 资源 del self.model_runner #等待所有子进程结束 for p in self.ps: p.join() #根据prompt创建sequence并递送给scheduler def add_request(self, prompt: str | list[int], sampling_params: SamplingParams): #支持输入字符串或者token_id因此加上判断 if isinstance(prompt, str): #tokenizer.encode/decode把输入文本转成 token_id或把生成的 token_id 转回文本 prompt self.tokenizer.encode(prompt) seq Sequence(prompt, sampling_params) self.scheduler.add(seq) #对sequence进行调度、前向推理和后处理 def step(self): seqs, is_prefill self.scheduler.schedule() #如果是prefill阶段返回填充的所有token如果是推理由于自回归一次出一个token因此一批次只会产生len(seqs)个token num_tokens sum(seq.num_scheduled_tokens for seq in seqs) if is_prefill else -len(seqs) #forward token_ids self.model_runner.call(run, seqs, is_prefill) #后处理 self.scheduler.postprocess(seqs, token_ids, is_prefill) #如果有sequence调度完成返回完成的prompt模型的回答 outputs [(seq.seq_id, seq.completion_token_ids) for seq in seqs if seq.is_finished] return outputs, num_tokens #如果既没有等待调度的sequence也没有在跑的就返回 def is_finished(self): return self.scheduler.is_finished() def generate( self, prompts: list[str] | list[list[int]], sampling_params: SamplingParams | list[SamplingParams], use_tqdm: bool True, ) - list[str]: #创建生成任务的进度条 pbar tqdm(totallen(prompts), descGenerating, dynamic_ncolsTrue, disablenot use_tqdm) #如果没有给每一个请求分别指定推理参数就复用主要是应对分布式请求 if not isinstance(sampling_params, list): sampling_params [sampling_params] * len(prompts) for prompt, sp in zip(prompts, sampling_params): self.add_request(prompt, sp) outputs {} #用于计算速度 prefill_throughput decode_throughput 0. while not self.is_finished(): #记录 step 开始时间后续计算prefill/decode吞吐量tok/s t perf_counter() output, num_tokens self.step() if num_tokens 0: prefill_throughput num_tokens / (perf_counter() - t) else: decode_throughput -num_tokens / (perf_counter() - t) #显示吞吐量等实时指标 pbar.set_postfix({ Prefill: f{int(prefill_throughput)}tok/s, Decode: f{int(decode_throughput)}tok/s, }) for seq_id, token_ids in output: outputs[seq_id] token_ids #每完成一个序列就更新进度 pbar.update(1) pbar.close() #保证输出的顺序和输入 prompts 的顺序一致 outputs [outputs[seq_id] for seq_id in sorted(outputs.keys())] outputs [{text: self.tokenizer.decode(token_ids), token_ids: token_ids} for token_ids in outputs] return outputs一些问题1.tqdm进度条参数解析 totallen(prompts) 进度条总数 请求数量每个请求完成时 update(1) descGenerating 进度条前面的标签文字 dynamic_ncolsTrue 动态调整列宽适应终端宽度 disablenot use_tqdm use_tqdmFalse 时禁用进度条不显示 3.如果step中的seq.isfinished不为Trueoutputs是不是一直为空还是说这里会等待直到变为True outputs会为空如果当前批次没有任何 seq 完成 代码不会等待seq.is_finished变为 Truestep是 “单次执行” 的函数调用一次就返回一次结果。 4.为什么step中的num_tokens在Decode 阶段统计 “生成的 Token 数”是负数 如果是正数代表当前这步做的是 Prefill提示词预填充处理的 token 数就是这个正数本身把所有 sequence 的长度加起来。 如果是负数代表当前这步做的是 Decode逐字生成。为什么用 -len(seqs)因为 Decode 阶段每个序列在一步里必定且只会生成 1 个 Token所以生成的总 Token 数就是当前序列的数量len(seqs)。给它加个负号外层一看到小于 0就知道“哦这是 Decode 阶段生成了 abs(num_tokens) 个 token”。 这在 generate 计算吞吐量时就能看出来decode_throughput -num_tokens / ...再次取负变成正数进行计算。 阶段 任务特性 token 数统计逻辑 数值正负的意义 Prefill预填充 处理用户输入的 prompt长文本一次性计算所有 prompt 的 token sum(len(seq) for seq in seqs)统计本次批次所有 prompt 的总 token 数正数 正数表示 “预填充阶段处理的总 token 数” Decode解码 逐 token 生成每个序列每次只生成 1 个 token批次内有 N 个序列就生成 N 个 token -len(seqs)用 “负的序列数” 表示 “解码阶段生成的 token 数”因为每个序列生成 1 个总 token 数 序列数负数是为了和 prefill 区分 负数的绝对值 解码阶段生成的 token 数 5.我没太看懂step的运行逻辑请你讲解 调度 (self.scheduler.schedule())问 Scheduler 拿一批当前可以执行的序列并确定这次是要算 Prefill 还是 Decode。 执行 (self.model_runner.call(...))把这批序列送进模型如果多卡则分发给多卡模型进行矩阵乘法等运算最后吐出一个包含新生成 Token ID 的列表。 后处理 (self.scheduler.postprocess(...))把新生成的 Token 塞到各自对应的 Sequence 屁股后面并且顺便检查一下有没有碰见 eos (结束符)或者有没有达到 max_tokens 长度限制如果碰到了就把它的状态标记为 FINISHED。 收获结果遍历刚才处理的序列把状态已经是 FINISHED 的序列挑选出来返回给外层。未完成的会在下一次 step 被继续处理。 6.self.tokenizer AutoTokenizer.from_pretrained(config.model, use_fastTrue)快速分词器和普通的分词器有什么区别 普通分词器use_fastFalse手动切菜 —— 纯 Python 实现逻辑直观但慢适合调试 / 定制 快速分词器use_fastTrue电动切菜机 —— 用 Rust/C 底层实现HuggingFace 的tokenizers库速度是普通版的 5~10 倍适合生产环境。 8.step中如果是prefill阶段也会执行postprocess吗postprocess 做什么 Prefill 阶段 1. 把 prompt 的 token_ids 追加到序列的completion_token_ids2. 几乎不会标记序列为 FINISHED因为 prefill 是处理输入还没开始生成不会碰到 eos 结束符3. 不会释放 KV Cache因为还要用 Cache 做后续生成。 Decode 阶段 1. 把生成的 1 个新 token 追加到序列2. 检查是否生成 eos 结束符 / 达到最大长度若是则标记序列为 FINISHED释放 KV Cache3. 从 running 队列移除完成的序列。 9.decode阶段每次调用一次step必定生成一个token吗不同序列都是一次step生成一个token吗速度是一致的吗 生成数量在 Decode 阶段每次调用一次 step每个正在运行的序列必定且只能生成 1 个 token。这是自回归大模型如 GPT、Qwen的物理规律决定的。 速度是否一致完全一致神同步。 在深度学习中这叫做批处理Batching。显卡不是一个序列一个序列地算而是把比如 16 个序列拼成一个大矩阵。显卡做一次矩阵乘法这 16 个序列的下一个 Token 就会在同一瞬间被同时算出来。所以它们是齐头并进的。 10.如果seqs中有序列完成了decode那么seqs执行step会怎么样 当一个序列比如它已经回答完毕生成了 eos 结束符在本次 step 中完成了 打上标记并清退在执行 postprocess 时系统发现它生成了结束符就会把它的状态改成 FINISHED并且立刻把它从 self.running正在运行的队列中一脚踢出去。 释放显存同时调用 self.block_manager.deallocate(seq)把这个句子之前占用的显存KV Cache全部还给系统。 上报给用户外层的 outputs 会捕捉到这个 FINISHED 状态把它最终的回答文本提取出来。 下一次 step在几毫秒后代码开启下一次 step()。当调度器再次调用 schedule() 去挑人时这个已经完成的序列已经不在运行队列里了所以它不会再参与下一次的运算。如果 waiting 队列里还有其他没开始的句子调度器这时候就有空位把新句子塞进来了。 12.pbar.set_postfix({ Prefill: f{int(prefill_throughputt)}tok/s, Decode: f{int(decode_throughput)}tok/s, })是不是意味着prefill或decode是根据本次step估算出的值是不是每次step都会更新这个值如果prefill结束prefill的速度值会消失还是维持最新的是不是所有prefill结束了才会decode还是说可以并行 prefill_throughput / decode_throughput 是每次 step 根据单次执行时间估算的不是滑动平均 每次 step 都更新所以值会持续变化 prefill 结束后维持最后一次计算的值不会消失代码里没有清零逻辑 先 prefill 完所有请求再开始 decode —— 不是并行的 具体看 scheduler.schedule() 的逻辑 优先处理 waiting 队列prefill 等 waiting 为空才处理 running 队列decode if scheduled_seqs: # prefill 还有就只做 prefill return scheduled_seqs, True # decode 只有当 waiting 空了才执行 while self.running and ...: ... return scheduled_seqs, False 但这里有个细节代码只会在 waiting 完全空时才进入 decode。如果有多个请求第一个请求 prefill 期间其他请求还在 waiting 队列等着——它们不会并行 prefill只会串行。scheduler.py作用解析作为调度核心被llm_engine.py调用这里最难的逻辑就是它跟BlockManger.py的交互这里会把kv_cache block分配给sequence让他记录到block_table简单理解就是#collections是 Python 内置的集合模块提供了比原生容器list/dict/tuple更高效、更专用的容器类型弥补了原生容器在特定场景下的性能或功能不足。常见的有 deque双端队列适用于高频头尾操作 defaultdict带默认值的字典 Counter计数字典 OrderedDict有序字典Python 3.7 dict 已有序但仍有专属方法等。 from collections import deque from nanovllm.config import Config from nanovllm.engine.sequence import Sequence, SequenceStatus from nanovllm.engine.block_manager import BlockManager class Scheduler: def __init__(self, config: Config): self.max_num_seqs config.max_num_seqs self.max_num_batched_tokens config.max_num_batched_tokens self.eos config.eos self.block_size config.kvcache_block_size self.block_manager BlockManager(config.num_kvcache_blocks, config.kvcache_block_size) self.waiting: deque[Sequence] deque() self.running: deque[Sequence] deque() def is_finished(self): #跟上个代码文件中说明一致必须是等待和运行队列均没有需要调度的sequence才会结束 return not self.waiting and not self.running def add(self, seq: Sequence): self.waiting.append(seq) #计算可以调度的sequence返回给llm_engine会发给model_runner进行前向计算 def schedule(self) - tuple[list[Sequence], bool]: scheduled_seqs [] num_batched_tokens 0 #可以看出必须所有prefill结束才能进入decode # prefill显然一批sequence只有队首Prefill完才到下一个 #需要注意这里是告诉llmengine调度哪些sequence不会给未命中的block计算kv,没有前向计算前向计算在llmengine的step的seqs, is_prefill self.scheduler.schedule() # ← ① 调度决定跑哪些 seq。token_ids self.model_runner.call(run, seqs, is_prefill) # ← ② 实际 forward 计算 while self.waiting and len(scheduled_seqs) self.max_num_seqs: seq self.waiting[0] #剩余token预算 remaining self.max_num_batched_tokens - num_batched_tokens if remaining 0: break #如果seq的kv_cache block为空 if not seq.block_table: #计算多少kv_cache block可以直接用 num_cached_blocks self.block_manager.can_allocate(seq) if num_cached_blocks -1: break num_tokens seq.num_tokens - num_cached_blocks * self.block_size else: num_tokens seq.num_tokens - seq.num_cached_tokens #如果空闲token不足只允许第一个sequence使用chunked_prefill if remaining num_tokens and scheduled_seqs: break if not seq.block_table: #把能直接用的kv_cache block记录到sequence的block_table这里是sequence的block_table第一次写入很重要这里涉及顺序理清以防止被绕晕 self.block_manager.allocate(seq, num_cached_blocks) seq.num_scheduled_tokens min(num_tokens, remaining) num_batched_tokens seq.num_scheduled_tokens if seq.num_cached_tokens seq.num_scheduled_tokens seq.num_tokens: seq.status SequenceStatus.RUNNING self.waiting.popleft() self.running.append(seq) scheduled_seqs.append(seq) if scheduled_seqs: return scheduled_seqs, True # decode显然一批sequence一个sequence一次生成一个token后就换下一个sequence生成 while self.running and len(scheduled_seqs) self.max_num_seqs: seq self.running.popleft() while not self.block_manager.can_append(seq): if self.running: self.preempt(self.running.pop()) else: self.preempt(seq) break else: seq.num_scheduled_tokens 1 seq.is_prefill False self.block_manager.may_append(seq) scheduled_seqs.append(seq) assert scheduled_seqs self.running.extendleft(reversed(scheduled_seqs)) return scheduled_seqs, False #将某个运行中的sequence剥夺资源变回等待态插入等待队列队首 def preempt(self, seq: Sequence): seq.status SequenceStatus.WAITING #被抢占的sequence需要重新prefill这个标志会用再model_runner、postprocess和sequence的getstate中 seq.is_prefill True #取消sequence对kv_cache块的引用如果这些kv_cache块不被其他sequence引用则将其放回空闲池只设置引用数为0不清空token和hash以便再次利用 self.block_manager.deallocate(seq) self.waiting.appendleft(seq) #后处理 def postprocess(self, seqs: list[Sequence], token_ids: list[int], is_prefill: bool): for seq, token_id in zip(seqs, token_ids): #当一个sequence被调度一次后将其被新分配的把Block的写入Block_manager管理器的记录表 self.block_manager.hash_blocks(seq) #prefill阶段可能num_scheduled_tokens会很多但decode阶段固定加1 seq.num_cached_tokens seq.num_scheduled_tokens seq.num_scheduled_tokens 0 if is_prefill and seq.num_cached_tokens seq.num_tokens: continue #prompt都处理完模型产生回答 seq.append_token(token_id) #如果回答结束或超出限制就强制结束 if (not seq.ignore_eos and token_id self.eos) or seq.num_completion_tokens seq.max_tokens: seq.status SequenceStatus.FINISHED self.block_manager.deallocate(seq) self.running.remove(seq)一些问题1.为什么如果有Prefill的序列就会先执行 优先级更高新请求Prefill是用户的首次请求需优先响应避免 Decode 阶段的 “老请求” 抢占全部资源 性能最优Prefill 计算量远大于 Decode需处理整个 prompt 的 KV Cache批量处理 Prefill 能最大化硬件利用率GPU 擅长批量并行计算 逻辑依赖Prefill 完成后才能进入 Decode 阶段若不优先处理 Prefill后续 Decode 也无法执行。 2.schedule函数返回的bool值是什么意思 该 bool 值标识本次调度的序列处于「Prefill 阶段」True还是「Decode 阶段」False核心作用是给 ModelRunner 传递阶段信息因为两个阶段的计算逻辑完全不同 True本次调度的是Prefill 序列新请求需计算整个 prompt 的 KV Cache注意力计算范围是全序列 False本次调度的是Decode 序列已有请求复用 KV Cache仅计算最后一个 token 的 KV注意力只关注最后一个位置。 3.为什么seq.block_table如果为空num_cached_blocks self.block_manager.can_allocate(seq)会不是0按理来说Block_table记录的是seq被分配的block吧如果不是的话seq最开始的block是哪来的。 Block物理块BlockManager.blocks 里的一个元素是真实存储 token 数据的内存块有一个 block_id。 block_table序列的块表一个序列的 block_table 是该序列分配到的物理 block_id 列表。 当 seq.block_table [] 时 该序列还没分配过任何物理块。但这不代表系统里没有可复用的块。 can_allocate 检查的是 hash 表hash_to_block_id里有没有相同前缀的物理块 h self.compute_hash(token_ids, h) block_id self.hash_to_block_id.get(h, -1) # 用 hash 查有没有现成的块 4.只能给第一个 seq 分块chunked prefill是什么意思。 看 postprocess 中的更新逻辑 seq.num_cached_tokens seq.num_scheduled_tokens # 累加已处理的 token 数 假设一个 seq 有 1000 个 prompt token第一次调度时 remaining200 num_scheduled_tokens 200本次实际处理 200 个 postprocess 后 num_cached_tokens 200 下一次 schedule() 调用时 else: num_tokens seq.num_tokens - seq.num_cached_tokens # 1000 - 200 800 这次只需要处理剩余的 800 个 token 了。如果预算够就继续处理直到完成。 chunked prefill 就是把一个大 seq 的 prefill 拆到多个 step 里做每次处理一部分。 5.这里似乎并没有对is_prefill的判断和设置is_prefill用在哪 model_runner.call(run, seqs, is_prefill) — 传给 ModelRunner告诉它按 prefill 还是 decode 方式执行 postprocess 中的条件判断scheduler.py 第 86-87 行 if is_prefill and seq.num_cached_tokens seq.num_tokens: continue # prefill 还没完成不追加 token seq.append_token(token_id) prefill 阶段还没处理完整个 seq 时不追加新 token因为新 token 还没生成。 Sequence.__getstate__sequence.py 第 73 行 last_state self.last_token if not self.is_prefill else self.token_ids 序列化时的区别decode 阶段只需保存最后一个 token之前的状态已在 KV cache 里prefill 阶段需保存完整 token_ids。本系列文章(待写完修正)[系列(1) 开篇](https://blog.csdn.net/xxx/article/details/xxxxxx)[系列(2) 核心原理](https://blog.csdn.net/xxx/article/details/xxxxxx)上一篇[系列(1)开篇](https://blog.csdn.net/xxx/article/details/xxxxxx)下一篇[系列(3)实战演练](https://blog.csdn.net/xxx/article/details/xxxxxx)