LMCache Request Stream SDK 深度解析:用有状态请求流编排 KV Cache 的多阶段推理
2026/9/15 14:01:58 网站建设 项目流程

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:

  1. prefill:把 prompt 的 KV 写入缓存;
  2. 修改缓存 KV:取回 KV、做压缩/裁剪等编辑、再写回;
  3. 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_tokensoutput_text),不受压缩等 KV 操作影响每次generate累加
done某次generate产出 token 数< max_tokens(视为触发 EOS)时为 True;update被调用时重置为 Falsegenerate判定;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、文本与时间间隔;若迭代中途抛异常,会包装为LMCacheRequestStreamErrorfinally中仍会累计已生成部分);
  • 返回本次调用的StreamPerfMetrics

modify_kv(fn, timeout=30.0, poll_interval=0.2)

编辑缓存 KV 的高层入口(request.py#L309-L339):

  1. 对每个注册的 context 调用retrieve(轮询等待缓存就绪);
  2. 从 KV 张量的 token 维tensors[KV].shape[2]得到cached_len
  3. tokens[cached_len:]连同既有_suffix_tokens存入_suffix_tokens
  4. 调用fn(tensors, tokens[:cached_len])得到编辑后的(new_kv, new_tokens),其中tensorsMapping[LMCacheSDKCacheKind, torch.Tensor],覆盖该算法用到的每种缓存类型;
  5. 通过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):

  • retrieveself.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/completionsstream=True),逐行解析data:事件,把每个 choice 包装为lmc_request.TokenEvent(token_id=..., text=...)产出;cache_salt非空时还会透传给 vLLM。也就是说,接入任意支持流式补全的引擎,只需提供一个同样签名的回调即可。

端到端实战:Token Dropping 工作流

examples/token_dropping/random_token_dropping.ipynb 给出了完整可运行的链路,核心步骤:

  1. 启动 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
  2. 启动 vLLM并接入LMCacheMPConnectorkv_connector_extra_config指向 6555 端口),注意加--return-tokens-as-token-ids以便逐 token 回调拿到 token id;

  3. 构造 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()
  4. 清空缓存后,用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通过线程池并发调用各流的generateget_perf_metrics(duration, fmt, width, mode, ...)聚合为Metricsmode"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_idstore-/retrieve-<uuid>)只标识请求会话,不参与缓存身份。
  • 传输:数据面优先走共享内存(SHM),否则回退 pickle,均由ContiguousTransferWrapper屏蔽差异——SDK 从不分支判断传输方式;这也解释了为什么示例中可以透明地依赖 SHM。
  • 注册前提:模型布局必须已由某个 vLLM 实例通过REGISTER_KV_CACHE注册到服务器,SDK 从/config/status读取chunk_sizekv_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_kvtimeout默认 30 秒、poll_interval默认 0.2 秒;超时抛LMCacheRequestStreamErrorgenerate中途异常也会包装为该错误并携带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_tokensoutput_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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询