LMCache Request Stream SDK 深度解析:用有状态请求流编排 KV Cache 的多阶段推理
【免费下载链接】LMCacheLMCache: Supercharge Your LLM with the Fastest KV Cache Layer项目地址: https://gitcode.com/GitHub_Trending/lm/LMCache
导读
LMCacheRequestStream是 LMCache SDK 提供的有状态、单一请求级编排层,它把一次完整推理请求(prefill → 修改缓存 KV → decode)绑定为一个逻辑流,让开发者无需手工在多个无状态LMCacheSDKContext调用之间搬运 token 状态。本文以 docs/design/sdk/request.md 为主线,结合 lmcache/sdk/request.py、lmcache/sdk/context.py 等源码与 token_dropping 示例,系统讲解其状态模型、Suffix 契约、完整公开 API、性能指标字段,以及如何用它实现「取 KV → 编辑 KV(如 token dropping)→ 放回 KV → 继续生成」的实战工作流。读完本文,你将掌握 Request Stream 的全部公开接口、内部状态流转语义与端到端接入方法。
为什么需要 Request Stream:无状态 Context 的痛点
按照 context.md 的定位,LMCacheSDKContext是无状态的:它的retrieve()/store()只按 token id 将 KV(以及可选的 query)张量移入移出正在运行的 LMCache MP 服务器,一次调用结束即"遗忘"。
但对于 token dropping 这类真实需求,一次请求会横跨多个推理 pass:
- prefill:把 prompt 的 KV 写入缓存;
- 修改缓存 KV:取回 KV、做压缩/裁剪等编辑、再写回;
- decode:基于编辑后的 KV 继续生成。
这三个 pass 在引擎侧会被编码成不同的请求,如果由用户自己维护"上一次缓存到哪里了、还有哪些 token 没进缓存",极易出错。LMCacheRequestStream正是为消除这份心智负担而生:它把属于同一个逻辑请求的所有推理 pass 绑定在一条流上,自动管理 token 序列、缓存对齐边界、已生成 token 计数与结束标志。
定位关系:RequestStream 本身不单独使用,LMCacheBatchedStream 会包装多条 RequestStream 用于批量提交请求;其底层依赖的
LMCacheSDKContext/LMCacheSDKCacheKind定义见 context.md。
一段代码看懂核心用法
原文档给出的最小示例完整展示了 Request Stream 的三段式生命周期:
import lmcache.sdk as lmc_sdk kind = lmc_sdk.LMCacheSDKCacheKind.KV ctx = lmc_sdk.connect(kind=kind, url=..., http_url=..., model_name=...) request = lmc_sdk.request.create_request( [ctx], post_completion, prompt_token_ids=source_tokens # contexts: iterable ) request.generate({"max_tokens": 1}) # prefill -> offload the prompt KV request.modify_kv(drop_tokens) # drop_tokens: retrieve -> edit -> store request.generate({"max_tokens": 256}) # replay the uncached tail + decode三个调用分别对应三个 pass:generate提交推理并触发 KV 卸载,modify_kv完成"取回-编辑-写回"整条链路,第二次generate重放未缓存的尾部并继续解码。
其中post_completion(prompt_token_ids, sampling_params, cache_salt)是注入的引擎调用函数(例如 vLLM 的/v1/completions流式接口),每生成一个 token 产出一次TokenEvent(token_id, text)。完整的可运行版本见 token-dropping 示例。
状态模型:流内每个字段的职责
在源码 lmcache/sdk/request.py 中,LMCacheRequestStream.__init__从初始 prompt 建立全部内部状态。整条逻辑序列是tokens与_suffix_tokens的组合,各字段语义如下:
| 字段 | 含义 | 何时变化 |
|---|---|---|
tokens | 支撑已存 KV 的 token 序列,作为下一次请求的 prompt 提交;KV 被update修改时整体替换 | generate追加生成 token;update替换为新序列 |
_suffix_tokens | 不在已存 KV 中的 token:modify_kv之后留下的非 chunk 对齐尾部,由下一次generate消费 | modify_kv记录;generate前置到 prompt 后清空 |
_decoded/_text_parts | 跨所有generate的累计生成 token 数 / 文本(对应decoded_tokens、output_text),不受压缩等 KV 操作影响 | 每次generate累加 |
done | 某次generate产出 token 数< max_tokens(视为触发 EOS)时为 True;update被调用时重置为 False | generate判定;update重置 |
request_stream_id在构造时生成(str(uuid.uuid4())),用于在批量场景中唯一标识一条流。
done 的判定细节
见 request.py#L227-L230:
# produces less than max_tokens --> EOS output_tokens = len(gen_tokens) max_tokens = sampling_params.get("max_tokens", 1) self.done = output_tokens < max_tokens当max_tokens未显式提供时默认为 1——这正是 prefill 阶段(强制max_tokens=1)能判定结束的原因。update()被调用后会重置done = False,因为缓存内容已改变,需要重新开始生成流程。
Suffix 契约:chunk 对齐边界如何被跨越
retrieve只返回chunk 对齐的前缀:LMCacheSDKContext.retrieve会把 token 数截断到(len(tokens) // chunk_size) * chunk_size(见 context.py#L400-L408),不足一个 chunk 时直接返回None。因此modify_kv编辑完 KV 后,剩余部分(sub-chunk 尾部 + 尚未卸载的 token)在缓存里没有对应的 KV,必须由流自身携带跨越这次编辑:
modify_kv在 KV 修改完成后,把未缓存尾部tokens[cached_len:]记录进_suffix_tokens;generate提交请求前,把_suffix_tokens(连同调用方传入的suffix_tokens)前置拼接到tokens,随后清空。
对应源码在 request.py#L199-L202:
pending = self._suffix_tokens + list(suffix_tokens) self._suffix_tokens = [] if pending: self.tokens.extend(pending)这样即使缓存只覆盖了 prompt 前 90% 的 token,剩余 10% 也不会丢失,会在下一次推理时作为 prompt 的一部分被重新提交、重新生成 KV。
公开 API 全景
构造:LMCacheRequestStream(...)与create_request(...)
两者等价(create_request只是工厂函数,见 request.py#L93-L115):
LMCacheRequestStream( contexts: Iterable[LMCacheSDKContext], post_completion: PostCompletion, prompt_token_ids: Sequence[int], cache_salt: str = "", )contexts是可迭代对象,如[kv_ctx]或[kv_ctx, q_ctx](KV + Query 双上下文场景);构造时按ctx.kind建立 kind → context 的映射。cache_salt是按用户隔离的盐(默认为空字符串),参与缓存寻址,避免不同用户共享同一前缀的 KV。- 使用
lmcache.sdk.connect(kind=...)可同时获得 KV 与 QUERY 两种 kind 的 context,分派逻辑见 lmcache/sdk/init.py。
generate(sampling_params, suffix_tokens=())→StreamPerfMetrics
运行一次推理 pass 并把结果追加进流历史(request.py#L183-L242):
- 先将
_suffix_tokens与调用方suffix_tokens合并追加到tokens; - 调用
self.post_completion(self.tokens, sampling_params, self.cache_salt)流式取 token; - 遍历
TokenEvent统计生成 token、文本与时间间隔;若迭代中途抛异常,会包装为LMCacheRequestStreamError(finally中仍会累计已生成部分); - 返回本次调用的
StreamPerfMetrics。
modify_kv(fn, timeout=30.0, poll_interval=0.2)
编辑缓存 KV 的高层入口(request.py#L309-L339):
- 对每个注册的 context 调用
retrieve(轮询等待缓存就绪); - 从 KV 张量的 token 维
tensors[KV].shape[2]得到cached_len; - 把
tokens[cached_len:]连同既有_suffix_tokens存入_suffix_tokens; - 调用
fn(tensors, tokens[:cached_len])得到编辑后的(new_kv, new_tokens),其中tensors是Mapping[LMCacheSDKCacheKind, torch.Tensor],覆盖该算法用到的每种缓存类型; - 通过
update写回新 KV。
ModifyFnType的类型签名定义在 context.py#L45-L48:
ModifyFnType = Callable[ [Mapping[LMCacheSDKCacheKind, torch.Tensor], Sequence[int]], tuple[torch.Tensor, Sequence[int]], ]retrieve(kind, timeout=30.0, poll_interval=0.2)与update(kind, kv, tokens)
两者是对 context 同名方法的薄封装(request.py#L244-L307):
retrieve按self.tokens轮询取缓存,timeout秒内拿不到即抛LMCacheRequestStreamError;返回的 KV 形状为[2, L, T, D],Q 形状为[1, L, T, D];update调用ctx.store(kv, tokens, cache_salt)写回编辑后的 KV,然后self.tokens = list(tokens)并self.done = False;若 store 报告该 KV 已缓存(去重命中),会记录一条 warning 而非报错。
访问器属性
request_stream_id:流的唯一 ID;suffix_tokens:待追加到 prompt 的尾部 token 列表;decoded_tokens:跨所有段累计生成的 token 数;output_text:跨所有段拼接的生成文本;output_tokens:当前完整 token 序列(含生成部分);is_done:是否已结束(EOS)。
StreamPerfMetrics:单次 generate 的性能报告
每次generate返回冻结的StreamPerfMetrics(dataclass,request.py#L35-L55),所有时间单位均为秒:
| 字段 | 含义 |
|---|---|
duration | 本次generate()调用耗时(秒) |
input_tokens | 输入 token 数(prompt + suffix) |
output_tokens | 本次调用生成的 token 数 |
input_tput | 输入吞吐(tokens/s) |
output_tput | 生成吞吐(tokens/s) |
tpot | 相邻生成 token 之间的时间间隔列表(秒),首对间隔被计作 TTFT |
ttft | 首 token 时间(秒) |
实现上,tpot记录time.perf_counter()的逐 token 差值,ttft取第一个间隔、tpot取其余间隔(request.py#L240-L241)。在批量场景中,LMCacheBatchedStream会把每条流的这份指标聚合成统一的Metrics报告(见 batch.md)。
PostCompletion 协议与真实实现
PostCompletion是定义引擎调用的Protocol(request.py#L71-L90):
def __call__( self, prompt_token_ids: list[int], sampling_params: dict[str, Any], cache_salt: str, ) -> Iterable[TokenEvent]: ...TokenEvent(request.py#L58-L69)只含两个字段:token_id(生成 token 的 id)与text(该 token 的解码文本)。
仓库中最直接的参考实现是 examples/token_dropping/utils.py#L258-L307 的make_post_completion:它用httpx.stream以 SSE 方式 POST vLLM 的/v1/completions(stream=True),逐行解析data:事件,把每个 choice 包装为lmc_request.TokenEvent(token_id=..., text=...)产出;cache_salt非空时还会透传给 vLLM。也就是说,接入任意支持流式补全的引擎,只需提供一个同样签名的回调即可。
端到端实战:Token Dropping 工作流
examples/token_dropping/random_token_dropping.ipynb 给出了完整可运行的链路,核心步骤:
启动 LMCache server(使用共享内存传输时指定
--shm-name与--no-l1-use-lazy):lmcache server \ --l1-size-gb 150 \ --eviction-policy LRU \ --chunk-size 256 \ --port 6555 \ --http-port 8080 \ --shm-name lmcache_kvcache_sdk_e2e \ --no-l1-use-lazy \ --supported-transfer-mode auto启动 vLLM并接入
LMCacheMPConnector(kv_connector_extra_config指向 6555 端口),注意加--return-tokens-as-token-ids以便逐 token 回调拿到 token id;构造 SDK context 与
post_completion,把多条请求加入LMCacheBatchedStream:batch = lmc_sdk.batch.LMCacheBatchedStream() for prompt in prompts: request = lmc_sdk.request.create_request( contexts=[ctx], post_completion=post_completion, prompt_token_ids=prompt, ) batch.add(request) results = batch.prefill( sampling_params={"max_tokens": 1, "temperature": 1.0, "ignore_eos": True} ) results.emit() results = batch.decode( sampling_params={"max_tokens": max_tokens, "temperature": 1.0, "ignore_eos": True} ) results.emit()清空缓存后,用
batch.modify(drop_tokens)对每条流的 KV 执行裁剪(modify_kv语义),再batch.decode对比吞吐。
批量层LMCacheBatchedStream的关键行为(batch.md 与 lmcache/sdk/batch.py):
prefill强制max_tokens=1,对全部流跑一遍generate并汇报 prefill 指标;modify(fn, ...)并发对每条流应用modify_kv,只汇报耗时;decode跑全部流并汇报 decode 指标;add/get_request_stream(stream_id)以request_stream_id为键管理成员;- 底层
run_request_streams通过线程池并发调用各流的generate,get_perf_metrics(duration, fmt, width, mode, ...)聚合为Metrics(mode取"prefill"或"decode"),支持终端表格emit()与to_dict()两种输出。
底层机制:chunk 对齐、寻址与传输
理解 Request Stream 的边界行为,需要知道它的存储语义来自 context.md:
- 缓存寻址:context 构建
IPCCacheServerKey(model_name, world_size=1, worker_id=0, token_ids, start=0, end=<chunk-aligned>, request_id, cache_salt);缓存身份 = token-chunk 哈希 +model_name+kv_rank(worker_id)+cache_salt,而request_id(store-/retrieve-<uuid>)只标识请求会话,不参与缓存身份。 - 传输:数据面优先走共享内存(SHM),否则回退 pickle,均由
ContiguousTransferWrapper屏蔽差异——SDK 从不分支判断传输方式;这也解释了为什么示例中可以透明地依赖 SHM。 - 注册前提:模型布局必须已由某个 vLLM 实例通过
REGISTER_KV_CACHE注册到服务器,SDK 从/config、/status读取chunk_size与kv_cache_layout并解码几何(HND 顺序、CPU 侧分配),无法仅凭model_name推导。 - 已知限制:当前仅支持
world_size == 1与单一非 hybrid kernel group;QUERY 类型在<model>##query键下由 vLLM worker 的 Q ring 通过REGISTER_Q_CACHE注册。
注意事项与边界行为
retrieve/update/modify_kv的timeout默认 30 秒、poll_interval默认 0.2 秒;超时抛LMCacheRequestStreamError,generate中途异常也会包装为该错误并携带request_stream_id。update会把流的 token 序列整体替换为new_tokens,因此编辑 KV 的算法必须返回与 KV 对应的完整 token 序列,且len(tokens)需与kv.shape[2]一致(store 侧会校验,见 context.py#L478-L486)。done依赖sampling_params["max_tokens"]的对比判定:prefill 阶段max_tokens=1是"必然产生结束信号"的约定,实际业务中若想禁止 EOS,示例使用了ignore_eos=True配合。output_tokens返回的是流内部tokens的引用,包含 prompt 与生成 token 的完整序列;如需纯生成部分,请结合decoded_tokens与output_text使用。
总而言之,LMCacheRequestStream把"无状态的 KV 移动原语"升级为"有状态的请求级编排",是 SDK 中实现 KV 编辑类优化(token dropping、KV 压缩、rerope 重算等)的核心抽象;配合LMCacheBatchedStream即可直接支撑离线批处理场景,其完整设计可继续阅读 context.md 与 batch.md,实战入口见 token_dropping 示例。
【免费下载链接】LMCacheLMCache: Supercharge Your LLM with the Fastest KV Cache Layer项目地址: https://gitcode.com/GitHub_Trending/lm/LMCache
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考