设计解析:EventGate 准入、Fencing 与广播分发)
LMCache 协调器缓存事件摄入层Cache-event Ingest设计解析EventGate 准入、Fencing 与广播分发【免费下载链接】LMCacheLMCache: Supercharge Your LLM with the Fastest KV Cache Layer项目地址: https://gitcode.com/GitHub_Trending/lm/LMCache本文是 LMCache 多进程协调器mp_coordinator缓存事件摄入层的深度技术指南以设计文档 docs/design/v1/mp_coordinator/ingest.md 为主体结合 ingest 模块源码 与单元测试展开。你将掌握协调器如何从单一事件流构建全集群缓存视图、EventGate 如何以Incarnation Fencing Seq 去重 Gap 检测三道机制保证准入有序、CacheEventBroadcaster 如何将事件扇出到 KeyDirectory 等消费方以及为何 L2 字节在进程重启后仍被保留而 L1 必须被 fencing。一、摄入层在整个协调器中的定位在 LMCache 的 mp_coordinator 架构中协调器关于整个 fleet集群内所有 MP server 实例缓存内容的全部认知都来自于一条事件流。摄入层ingest layer就是这条流进入协调器的唯一入口它决定两件事什么被准入admission——由EventGate负责谁能看到它fan-out——由CacheEventBroadcaster负责。设计上刻意让这一层不持有任何缓存状态状态全部由下游的消费方consumers持有。这正是摄入层与视图层KeyDirectory解耦的核心理念。数据流如下source adapter ingest layer consumers ────────────────────────────────────────────────────────────────────────────── POST /events ──▶ HttpCacheEventSource ──▶ EventGate ──▶ CacheEventBroadcaster (HTTP push) .ingest(batches) fence / .broadcast(batch) ────▶ KeyDirectory dedup / .fence_instance(id) ──▶ FleetEvictionController gap detect对应到仓库源码该层由四个模块构成lmcache/v1/mp_coordinator/ingest/模块职责event_source.py定义 source 的生命周期/状态契约CacheEventSourceProtocol、EventReplayCapability、CacheEventSourceStatushttp_event_source.py非持久化non-durable的POST /events推送源实现event_gate.pyEventGate准入控制fencing、去重、gap 检测event_broadcaster.pyCacheEventBroadcaster扇出器 CacheEventConsumer消费方协议契约词汇表定义在 lmcache/v1/mp_coordinator/api.py核心数据结构CacheEventBatchHTTP 表面是 http_apis/events_api.pyPOST /events。二、为什么要把 Gate 与 KeyDirectory 分开这是摄入层设计最核心的决策。文档给出的理由是这两件事的生命周期不同归属也不同。流的准入stream admission是**发射方emitter**的属性——每个 emitter 一个游标cursor无论状态是否被保留都有效而放置关系placements、使用量usage、LRU 是**缓存cache**的属性。如果把准入逻辑折叠进 KeyDirectory会产生两个问题KeyDirectory 会变成每个事件的强制第一个消费方——新增第二个消费方就必须绕道经过它回答这批 batch 是不是重放必须去问 placement store放置存储。将 Gate 独立出来后KeyDirectory 只是另一个普通的CacheEventConsumer按约定第一个注册作为事实来源 source of truth新增一个消费方只需在 create_app 装配函数 中调用一次register_consumer摄入层本身零改动。从 app.py 的实际装配代码可以看到这一约定event_broadcaster CacheEventBroadcaster() # Views first: a controller acts on the batch a view has consumed. for collaborator in (*views.all(), *controllers.all()): if isinstance(collaborator, CacheEventConsumer): event_broadcaster.register_consumer(collaborator) quiesce QuiesceLock() event_gate EventGate(event_broadcaster, quiesce) event_source HttpCacheEventSource(event_gate)EventGate是唯一被显式命名的持久化组件name stream_cursors因为它既不属 view 也不属 controller却要随 checkpoint 一起持久化详见下文游标持久化一节。三、批次Batch的数据结构在深入准入逻辑前先看摄入层的输入——CacheEventBatch。它定义于 lmcache/v1/mp_coordinator/api.py一个批次 (instance_id, incarnation, seq, event_type, tier, backend, entries[], shared, ts, ...)字段类型/取值语义instance_idstr非空发射方的唯一 ID本质上是emitter stream idincarnationint≥ 0发射方的重启计数器更大的值会 fencing 掉同一instance_id低版本上报的所有放置seqint≥ 1按(instance_id, incarnation)维度的单调批次计数event_typeCacheEventType本批次内所有条目发生了什么如STORE、DELETE、ACCESS、CONFIGtierTierl1/l2绝不等于all事件作用于的缓存层级backendstr如dram、cxl、fs、valkey层级内的存储后端对store/delete必填是放置身份的一部分对access为空entrieslist[CacheEventEntry]受影响的 key 列表config批次恒为空sharedbool默认 False后端是否为多个实例共享的存储域如同一个 S3 bucket 或 CXL 池tsfloat默认 0.0发射方墙钟秒数capacity_bytes/capacity_revisionint仅config声明 compartment 容量及其修订号CacheEventBatch.__post_init__通过ValueError强制内在不变量instance_id非空、seq 1、incarnation 0、tier必须是具体层级l1/l2而非all、config批次不得携带 entries、携带放置身份的批次backend不得为空。四、准入机制EventGate.ingestEventGate是摄入层的准入权威实现于 lmcache/v1/mp_coordinator/ingest/event_gate.py。它按顺序执行三道校验任一关卡失败即丢弃批次机制规则原因Incarnation fencingincarnation 当前值 → 丢弃批次STALE_INCARNATIONincarnation 当前值 → 对所有消费方执行fence_instance(id)然后开启全新游标重启会清空上报方的内存——它此前上报的 L1 放置不得继续存活。L2 字节在磁盘上跨重启持久因此 L2刻意不被 fencing只追踪 L2 的消费方对该钩子 no-opSeq 去重seq 已准入的最后一条同一 incarnation→ 丢弃批次DUPLICATE重放重试、事件总线重新投递必须是幂等的Gap 检测seq 已准入的最后一条 1→ 设置该 emitter 的gap_detected标志仍然准入事件可能丢失标志将该 emitter 的切片标记为过期直到流被重放依赖持久传输的 retention。消费方应用是幂等的所以跨 gap 准入是安全的对应 event_gate.py 核心实现with self._quiesce.applying(), self._lock: cursor self._cursors.get(batch.instance_id) if cursor is not None: if batch.incarnation cursor.incarnation: return IngestResult.STALE_INCARNATION if batch.incarnation cursor.incarnation: # Restart: the emitters memory is empty, so the L1 # facts its previous incarnation reported are void. self._broadcaster.fence_instance(batch.instance_id) cursor None elif batch.seq cursor.last_seq: return IngestResult.DUPLICATE if cursor is None: cursor _StreamCursor(incarnationbatch.incarnation) self._cursors[batch.instance_id] cursor if batch.seq cursor.last_seq 1 and not cursor.gap_detected: cursor.gap_detected True logger.warning(...) cursor.last_seq batch.seq self._broadcaster.broadcast(batch) return IngestResult.ADMITTED几个值得注意的实现细节准入结果的三种枚举IngestResult为ADMITTED/DUPLICATE/STALE_INCARNATION只有ADMITTED会到达消费方。按实例 FIFOseq是设计所需的唯一排序每个实例是自己事实的唯一写入者因此不存在全局顺序也不需要跨实例仲裁。instance_id本质上是emitter stream id——共享介质shared pool的 controller 用自己的稳定 ID 上报得到的只是一条普通去重、fenced 的流没有任何特判。锁跨越整个扇出过程EventGate的锁在 fan-out 全程持有with self._quiesce.applying(), self._lock:保证一个 emitter 的批次按准入顺序到达消费方。因此消费方绝不能回调进 Gate否则死锁这是消费方实现必须遵守的纪律。QuiesceLock的配合quiesce在每次变更调用期间持有使捕获持久状态的组件永远不会看到一个半应用的批次即捕获方与 ingest 方按相同顺序取两把锁。批量入口ingest_batchessource 适配器把有序批次列表喂给EventGate.ingest_batchesevent_gate.py#L139-L168它按列表顺序逐个调用ingest并汇总统计CacheEventIngestSummary(applied, duplicates, stale)三组计数。这是 HTTP API 响应的直接数据来源。drop_instance无批次的 fencingdrop_instance(id)event_gate.py#L170-L181是不带批次的 fencing——用于注销deregistration和心跳超时驱逐。它还会遗忘游标因此重连后可以从任意 incarnation 全新开始。文档明确标注把它接入注册表registry仍是后续工作但方法已存在且有测试覆盖。游标持久化capture / restore游标随 placement 一起走 checkpointPersistenceType.CHECKPOINT因为游标描述的是放置的来源流。capture()返回{cursors: {instance_id: (incarnation, last_seq, gap_detected)}}restore()要求 Gate 尚未准入过任何批次否则抛ValueError。event_gate.py 的注释 点明了为什么必须持久化游标没有游标的 Gate 无法 fencing——fencing 需要与先前的 incarnation 比较而空 Gate 没有可比较的对象重启后的 server 的陈旧 L1 切片会被永远公告出去。这一点与测试 tests/v1/mp_coordinator/persistence/test_restore_without_replay.py 的意图相互印证。五、Sources摄入是唯一入口ingest是唯一入口每个 source 都必须携带一个流(instance_id, incarnation, seq)。source 适配器把有序批次列表喂给EventGate.ingest_batches并上报汇总后的 admitted / duplicate / stale 计数。当前实现HttpCacheEventSource目前唯一的适配器是 http_event_source.py 中的HttpCacheEventSource由 MP-server 端的CacheEventSubscriber通过POST /events推送参见 cache_events.md。它是一个非持久化推送源FastAPI 拥有其请求生命周期start()/stop()均为空实现因为没有任何 source 自有的后台资源无法 seek 或重放那些在协调器接受之前就失败的事件。从 http_apis/events_api.py 可以看到POST /events表面非常薄把请求体中的body.batches交给HttpCacheEventSource.ingest返回CacheEventsResponse(applied, duplicates, stale)。重复批次与陈旧 incarnation 被丢弃并计数而非报错——这是幂等语义的直接体现。该端点被设计为顶层而非挂在/directory下因为这条流喂养的是所有消费方而不是其中某一个。生命周期契约刻意更小source 生命周期/状态契约event_source.py刻意比可重放源契约更小EventReplayCapability只有两个值NONE与SEEKABLEHTTP 源报告replay_capabilityNONE且不实现静默的 no-opseek未来的持久化 source 会在其位置有具体表示时才增加自己的传输位置与 seek/lag 契约。注意传输位置例如 Kafka 分区偏移量与 Gate 内部按 emitter 维护的 seq 游标是两套独立的概念。为何刻意没有扫描型入口一个扫描scan当前内容而非流式的 source——例如曾经用GET /cache/objects分页实现的启动期 L2 重新同步resync——没有流位置准入它需要一个第二个、无游标的入口。这个入口被刻意省略只有一个入口点就意味着只有一处地方决定排序与 fencing。文档给出的权衡建议是重新引入扫描源 重新引入reconcile并回答其批次对游标的影响而来自持久传输的事件重放则不需要任何新入口。这两条路要对照着评估。六、扇出CacheEventBroadcaster与消费方协议CacheEventBroadcaster实现于 event_broadcaster.py自身不持锁只要每个消费方各自线程安全fan-out 就是线程安全的。它维护一个有序的消费方列表broadcast(batch)按注册顺序逐个调用consumer.consume(batch)。消费方实现两个钩子CacheEventConsumerProtocolevent_broadcaster.py#L18-L44钩子语义实现要点consume(batch)应用一个已准入的批次按准入顺序调用每个事件每次投递尝试最多到达一次。跳过无关的 tier 与事件类型是消费方自己的职责fence_instance(instance_id)丢弃该实例在自己内存中持有的状态仅限 L1KeyDirectory丢弃该实例上报的 L1 放置借助其 per-instance 反向索引代价与该实例的 key 数成正比而非全表扫描FleetEvictionControllerno-op因为它核算的 L2 字节比进程更长寿只能通过DELETE离开注册顺序即调用顺序。当前的注册顺序是先 KeyDirectory放置与 token 绑定事实来源后 eviction controller按 salt 统计使用量与 LRU。两者相互独立——controller 自身的读后写read-after-write顺序是其内部属性参见 usage_and_eviction.md与注册顺序无关。七、状态在哪里三类问题的责任划分摄入层设计的一个核心问题是状态放在哪。文档用一张表明确划分了三类问题的归属问题向谁询问这个批次是重放吗这个 emitter 在哪个 incarnation我们丢事件了吗EventGate.stats()这个 key 存在哪里这个 chunk 持有哪些 tokenKeyDirectory这个 salt 用了多少字节该驱逐什么FleetEvictionController一个值得注意的限制EventGate.stats()目前没有 HTTP 端点——GET /directory/stats刻意只报告目录内容。因此gap_detected当前对运维人员不可见暴露它是重放集成后续工作的一部分。八、源码级验证测试如何印证三道机制单元测试 tests/v1/mp_coordinator/test_event_gate.py278 行用_RecordingConsumer忠实记录了 Gate 对消费方的每次调用逐条印证了设计文档的每个断言Seq 去重test_duplicate_seq_is_dropped_before_the_consumersseq1 已准入后再次提交 seq1即使 size 不同返回DUPLICATE消费方只收到 1 个批次test_replayed_older_seq_is_droppedseq2 之后重放 seq1 同样被丢弃消费方仅收到 2 个批次。Gap 检测test_seq_gap_sets_the_gap_flag_but_admitsseq1 后提交 seq5 →ADMITTEDstats()[node-a]显示gap_detectedTrue且last_seq5消费方收到 2 个批次跨 gap 照常应用test_contiguous_seqs_do_not_flag_gap连续 seq 不置标志。Incarnation fencingtest_new_incarnation_fences_consumers_before_admittingincarnation1 的 seq1 之后提交 incarnation2 的 seq1 →ADMITTED消费方收到fenced [node-a]游标更新为 incarnation2test_new_incarnation_drops_the_directory_l1_placements用真实KeyDirectory验证——incarnation1 上报 key(1)、key(2) 后incarnation2 上报 key(3) 会使 key(1)、key(2) 的 lookup 变为空列表而 key(3) 的 placement incarnation 为 2test_fence_spares_other_instances_placementsnode-a 重启 fencing 后不影响node-b 对同一 key 的放置——这正是per-instance 反向索引、按实例精确清理的佐证test_stale_incarnation_batch_is_droppedincarnation2 后提交 incarnation1即使 seq99→STALE_INCARNATION消费方无新增调用、无 fencingtest_same_incarnation_never_fences同一 incarnation 内连续批次绝不触发 fencing。drop_instancetest_drop_instance_fences_consumers_and_forgets_the_cursordrop 后stats()为空且后续以任意 incarnation 重连都能被正常准入test_drop_unknown_instance_is_noop_for_the_cursor对未知实例 drop 是安全 no-op。此外tests/v1/mp_coordinator/test_http_event_source.py、test_event_broadcaster.py 分别覆盖 source 与扇出器tests/v1/mp_coordinator/test_cache_events.py 覆盖端到端的事件上报路径。九、边界与后续工作deliberately out of scope文档明确列出的三项后续工作follow-ups也是当前版本刻意不做的事重放集成Replay integration通过 HTTP 暴露gap_detected然后借助持久传输的 retention 重放该 emitter 的流来消除 gap 标志。注册表集成Registry integration在注销/心跳超时驱逐时调用EventGate.drop_instance。共享池的分配代数Allocation generations为 shared pools 提供确定性的跨上报方冲突解决allocation generations。这些边界定义保证了摄入层的核心Gate Broadcaster 单入口保持最小且稳定——未来无论接入 Kafka 等持久消息队列、还是引入新的消费方fencing、去重、扇出与既有消费方都无需改动这正是把准入与状态分离这一设计决策的长远价值所在。速查摄入层关键路径一览关注点位置设计文档docs/design/v1/mp_coordinator/ingest.md批次数据结构CacheEventBatchlmcache/v1/mp_coordinator/api.py准入 Gatelmcache/v1/mp_coordinator/ingest/event_gate.py扇出与消费方协议lmcache/v1/mp_coordinator/ingest/event_broadcaster.pySource 契约lmcache/v1/mp_coordinator/ingest/event_source.pyHTTP source 实现lmcache/v1/mp_coordinator/ingest/http_event_source.pyPOST /events端点lmcache/v1/mp_coordinator/http_apis/events_api.py装配注册消费方、构建 Gatelmcache/v1/mp_coordinator/app.pyGate 单元测试tests/v1/mp_coordinator/test_event_gate.py关联设计事件上报docs/design/v1/mp_coordinator/cache_events.md关联设计KeyDirectorydocs/design/v1/mp_coordinator/views/key_directory.md关联设计使用量与驱逐docs/design/v1/mp_coordinator/usage_and_eviction.md【免费下载链接】LMCacheLMCache: Supercharge Your LLM with the Fastest KV Cache Layer项目地址: https://gitcode.com/GitHub_Trending/lm/LMCache创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考