Skip to content
Open
Show file tree
Hide file tree
Changes from 13 commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,8 @@ flowchart TD
py_connector["py_connector"] --> client["client SDK"]
client -.-> config
optimizer["optimizer(仿真与优化)"] -. cache_location .-> meta
optimizer --> protocol
optimizer -. SubscribeEvents .-> service
```

> 若模块职责、依赖方向或调用关系变动,或新增/删除模块,请同步更新本图与 [docs/design/module_architecture.md](docs/design/module_architecture.md)。
Expand Down
11 changes: 10 additions & 1 deletion docs/design/module_architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,7 @@ service → manager → meta → config → data_storage → common
### 3.3 客户端与 Optimizer

- **client** 是独立的对外分支,仅共享 `common`、`config`、`data_storage`、`protocol` 以及 `service/util:manager_message_proto_util`;**py_connector** 通过 pybind 位于 client 之上。核心服务端不依赖 client。运行时,元数据面经 gRPC(C++ `MetaClient`)或 HTTP(py_connector 的 Python `KvCacheManagerClient`)调用 KVCM 服务,数据面经 C++ `TransferClient` 直接读写存储后端——这几条链路是理解端到端流程的关键(见第 4 节)。
- **optimizer** 负责 KVCache 访问 trace 的仿真与优化(命中率/容量分析、逐出与容量参数调优),并通过独立 online runtime/service 提供实时 TraceQuery。full-attention LRU 的在线多容量统计复用 LiteHit;它通过 `meta:cache_location` 类型与 `event` 的 optimizer 事件与核心关联
- **optimizer** 负责 KVCache 访问 trace 的仿真与优化(命中率/容量分析、逐出与容量参数调优),并通过独立 online runtime/service 提供实时 TraceQuery。full-attention LRU 的在线多容量统计复用 LiteHit;在线进程还可通过通用服务发现找到全部 KVCM endpoint,在 KVCM 的 Meta gRPC 端口调用 `SubscribeEvents` 并直接回放事件。KVCM 不反向发现 Optimizer

### 3.4 模块关系图

Expand Down Expand Up @@ -146,6 +146,7 @@ flowchart TD
service --> data_storage
service --> protocol
manager --> protocol
event --> protocol
config --> protocol
meta --> config
meta --> data_storage
Expand All @@ -165,6 +166,8 @@ flowchart TD
%% optimizer 的关联
optimizer -. cache_location 类型 .-> meta
optimizer --> common
optimizer --> protocol
optimizer -. gRPC SubscribeEvents(运行时) .-> service

%% 底层通用模块被广泛依赖
meta --> common
Expand Down Expand Up @@ -307,6 +310,12 @@ flowchart LR

`LeaderElector`(config)基于 `CoordinationBackend`(memory/file/redis)的分布式锁选主。`Server` 在成为 Leader 时调用 `CacheManager::DoRecover` 恢复状态,随后启动 GC、恢复 Reclaimer、启动 MigrationManager 并开放 leader-only 请求。降级时先通知 GC/Reclaimer 停止新工作并关闭、排空 leader-only 请求,再 join GC、停止 MigrationManager,最后调用 `DoCleanup` 清理运行时状态(正在进行的写入按失败处理)。Python `KvCacheManagerClient` 使用服务发现 URL 时,会在每次 Leader 刷新前重新选择一个 Manager 发现端点,避免把 Leader 查询入口固定在单个节点上。

### 4.9 KVCM 到在线 Optimizer 的事件流

`CacheManager` 在读取路径产生 `CacheGetEvent`,`EventManager` 将其同时交给各 publisher。启用 optimizer publisher 后,单 worker 按事件顺序转换为 `TraceQueryRequest`,再写入每个 gRPC 订阅者的独立有界队列。KVCM 的 `OptimizerEventStreamService` 注册在现有 Meta gRPC server 上,不新增端口。

Optimizer 作为客户端通过 `ServiceDiscoveryFactory` 获取 KVCM seed endpoint,再调用 `MetaService.GetClusterInfo` 定位当前 Leader;任意时刻只向 Leader 保持一条 `SubscribeEvents` response stream。supervisor 周期刷新服务发现和 Leader,并通过同一 Meta gRPC 端口上的 `OptimizerEventStreamService.GetConfiguration` 拉取 Instance Group / Instance 快照,按 Group 后 Instance 的顺序自动注册新增配置;未知 `instance_id` 会立即触发一次额外刷新。切主时先同步配置,再迁移 stream。stream 收到事件后直接调用 `OnlineOptimizerManager::TraceQuery`,Optimizer 侧不再增加业务队列。

---

## 5. 修改功能前的提示
Expand Down
3 changes: 3 additions & 0 deletions kv_cache_manager/event/optimizer_event_publisher.cc
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,9 @@ bool OptimizerEventPublisher::Convert(const std::shared_ptr<BaseEvent> &event,
for (const auto key : get_event->get_keys()) {
out->add_block_keys(key);
}
for (const auto &location_spec_name : get_event->location_spec_names()) {
out->add_location_spec_names(location_spec_name);
}
// The event carries microseconds; the replay works in nanoseconds.
out->set_timestamp_ns(get_event->event_trigger_time_us() * 1000);

Expand Down
16 changes: 14 additions & 2 deletions kv_cache_manager/event/test/optimizer_event_publisher_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -72,11 +72,12 @@ std::shared_ptr<CacheGetEvent> MakeGetEvent(const std::string &instance_id,
const std::vector<std::int64_t> &keys,
const std::vector<std::int64_t> &tokens,
const std::string &query_type = "prefix_match",
const std::string &trace_id = "trace-1") {
const std::string &trace_id = "trace-1",
const std::vector<std::string> &location_spec_names = {}) {
auto event = std::make_shared<CacheGetEvent>(instance_id);
event->SetEventTriggerTime();
event->set_trace_id(trace_id);
event->SetAddtionalArgs(query_type, keys, tokens, BlockMask(), 0, {});
event->SetAddtionalArgs(query_type, keys, tokens, BlockMask(), 0, location_spec_names);
return event;
}

Expand Down Expand Up @@ -239,6 +240,17 @@ TEST_F(OptimizerEventPublisherTest, TestTraceIdIsCarriedThrough) {
EXPECT_EQ("trace-abc", request.trace_id());
}

TEST_F(OptimizerEventPublisherTest, TestLocationSpecNamesAreCarriedThrough) {
ASSERT_TRUE(publisher_->Init(""));
ASSERT_TRUE(publisher_->Publish(
MakeGetEvent("instance-a", {1, 2}, {1, 2, 3}, "prefix_match", "trace-abc", {"tp0", "tp1"})));
ASSERT_TRUE(WaitForEvents(*sink_, 1));

const auto request = sink_->Events()[0];
EXPECT_EQ((std::vector<std::string>{"tp0", "tp1"}),
std::vector<std::string>(request.location_spec_names().begin(), request.location_spec_names().end()));
}

TEST_F(OptimizerEventPublisherTest, TestDoesNotFilterQueryType) {
ASSERT_TRUE(publisher_->Init(""));
ASSERT_TRUE(publisher_->Publish(MakeGetEvent("instance-a", {1, 2}, {1, 2}, "reverse_roll_sw_match")));
Expand Down
71 changes: 31 additions & 40 deletions kv_cache_manager/optimizer/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -48,39 +48,6 @@ cc_library(
],
)

cc_library(
name = "optimizer_metrics_collector",
srcs = [
"service/metrics/optimizer_metrics_collector.cc",
],
hdrs = [
"service/metrics/optimizer_metrics_collector.h",
],
deps = [
"//kv_cache_manager/common",
"//kv_cache_manager/metrics:metrics_registry",
],
)

cc_library(
name = "optimizer_metrics_reporter",
srcs = [
"service/metrics/optimizer_metrics_reporter.cc",
],
hdrs = [
"service/metrics/optimizer_metrics_reporter.h",
],
deps = [
"//kv_cache_manager/optimizer/manager:online_optimizer_manager",
":optimizer_metrics_collector",
"//kv_cache_manager/common",
"//kv_cache_manager/common:common_util",
"//kv_cache_manager/metrics:kmon_param",
"//kv_cache_manager/metrics:metrics_registry",
"@havenask//aios/kmonitor:kmonitor_client_cpp",
],
)

cc_library(
name = "optimizer_call_guard",
srcs = [
Expand All @@ -90,8 +57,8 @@ cc_library(
"service/optimizer_call_guard.h",
],
deps = [
":optimizer_metrics_collector",
":optimizer_metrics_reporter",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_collector",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_reporter",
"//kv_cache_manager/common",
"@rapidjson",
],
Expand All @@ -108,8 +75,8 @@ cc_library(
deps = [
"//kv_cache_manager/optimizer/manager:online_optimizer_manager",
":optimizer_call_guard",
":optimizer_metrics_collector",
":optimizer_metrics_reporter",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_collector",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_reporter",
"//kv_cache_manager/common",
"//kv_cache_manager/optimizer/config:online_optimizer_config",
"//kv_cache_manager/optimizer/config:optimizer_registry_manager",
Expand All @@ -126,7 +93,7 @@ cc_library(
"service/grpc/optimizer_service_grpc.h",
],
deps = [
":optimizer_metrics_collector",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_collector",
":optimizer_service_impl",
"//kv_cache_manager/common",
"//kv_cache_manager/metrics:metrics_registry",
Expand All @@ -149,7 +116,7 @@ cc_library(
"-fcoroutines",
],
deps = [
":optimizer_metrics_collector",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_collector",
":optimizer_service_impl",
"//kv_cache_manager/common",
"//kv_cache_manager/metrics:metrics_registry",
Expand All @@ -174,6 +141,28 @@ cc_library(
],
)

cc_library(
name = "kvcm_event_subscriber",
srcs = [
"service/event_subscriber/kvcm_event_subscriber.cc",
],
hdrs = [
"service/event_subscriber/kvcm_event_subscriber.h",
],
deps = [
":online_optimizer_server_config",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_collector",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_reporter",
":optimizer_service_impl",
"//kv_cache_manager/common:logger",
"//kv_cache_manager/common:service_discovery",
"//kv_cache_manager/common:service_discovery_factory",
"//kv_cache_manager/protocol/protobuf:service_cc_grpc",
"//kv_cache_manager/service/util:common",
"@grpc//:grpc++",
],
)

cc_library(
name = "online_optimizer_server",
srcs = [
Expand All @@ -188,12 +177,14 @@ cc_library(
],
deps = [
"//kv_cache_manager/optimizer/manager:online_optimizer_manager",
":kvcm_event_subscriber",
":online_optimizer_server_config",
":optimizer_metrics_reporter",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_reporter",
":optimizer_service_grpc",
":optimizer_service_http",
":optimizer_service_impl",
"//kv_cache_manager/common",
"//kv_cache_manager/common:loop_thread",
"//kv_cache_manager/metrics:metrics_registry",
"//kv_cache_manager/optimizer/config:optimizer_registry_manager",
"@grpc//:grpc++",
Expand Down
33 changes: 33 additions & 0 deletions kv_cache_manager/optimizer/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,39 @@ OptimizerManager (core coordinator)
HitAnalysis (result analysis)
```

### 在线订阅 KVCM 事件

`online_optimizer_server_main` 可以主动发现 KVCM,并通过 KVCM 现有的 Meta gRPC 端口订阅缓存读取事件。KVCM 不需要知道 Optimizer 地址,也不新增事件端口;Optimizer 是 gRPC 客户端,调用 `OptimizerEventStreamService.SubscribeEvents`,KVCM 通过 response stream 写入 `TraceQueryRequest`。

KVCM 侧先在原有 server 配置中启用 optimizer publisher:

```properties
kvcm.event.event_publishers_configs={"log":{"enable":true,"queue_size":10000},"optimizer":{"enable":true,"queue_size":100000,"max_subscribers":4,"subscriber_queue_size":10000}}
```

Optimizer 侧在现有 JSON 配置中加入订阅配置,不使用额外配置文件:

```json
{
"kvcm_event_subscription": {
"enable": true,
"service_discovery_url": "static://127.0.0.1:6381",
"consumer_id": "online-optimizer",
"discovery_refresh_interval_ms": 5000
}
}
```

`service_discovery_url` 指向 KVCM 的 `kvcm.service.rpc_port`,支持通用服务发现 URL(如 `static://`、`vipserver://`、`spectrum://`)。发现结果只作为 seed:Optimizer 周期调用任一健康 seed 的 `MetaService.GetClusterInfo` 获取当前 Leader,并且只向 Leader 维持一条 `SubscribeEvents` stream。切主后下一次刷新会先同步新 Leader 的配置,再关闭旧 stream、连接新 Leader;断流期间会自动重连。

`OptimizerEventStreamService.GetConfiguration` 与事件流共用 KVCM 的 Meta gRPC 端口。Optimizer 启动时及每次服务发现刷新时拉取一次 Instance Group / Instance 快照,先创建缺失的 Group,再注册缺失的 Instance;收到未知 `instance_id` 时还会立即唤醒一次配置刷新。因此 KVCM 新增实例后不需要再提前调用 Optimizer 的注册 API。当前同步只添加新配置,不删除或热更新已经存在的 Optimizer 配置。

自动注册采用能从 KVCM 配置直接确定的口径:Instance Group 的 quota byte 数转换成一个 GiB 容量点,使用当前 Optimizer 支持的 LRU,并开启 prefix hash;Instance 按 full-only 注册,优先合并名称以 `full` / `FULL` 开头的 location spec group,没有时使用全部 spec。当前在线 indexer 不支持 shared group quota,因此同组各 Instance 分别按完整 Group quota 模拟。KVCM 没有 `linear_step` 等 Optimizer 专属参数,因此这里不猜测 linear/mamba 周期。

订阅器固定使用两个线程:一个 supervisor 线程负责服务发现、Leader 查询和配置同步,一个 stream 线程负责读取当前 Leader 的事件;收到事件后直接调用 `OnlineOptimizerManager`,不增加额外事件队列。未知 Instance 的首条事件会记录并丢弃,配置刷新完成后的后续事件正常进入统计。事件时间戳用于 LiteHit 和线性 indexer 的 TTL 判定,旧客户端未设置时间戳时仍回退到 Optimizer 本机墙钟。

在线 full-attention 实例还会输出 `mrc` gauge(Prometheus 名称默认为 `kvcm_optimizer_mrc`,标签为 `instance_id`,单位 byte)。它表示最近一个 `metrics_report_interval_ms` 上报周期内,达到理论无限容量命中数 95% 所需的最小 LRU 容量;每次上报会原子取走并清空仅供 MRC 使用的 hit curve,不影响查询数、命中率等累计指标。该值直接聚合 LiteHit 产生的容量无关 hit curve,不依赖预先配置的离散容量点;周期内尚无理论可命中 block 时值为 0。


### Eviction Policies

Expand Down
Loading
Loading