Skip to content
Open
Show file tree
Hide file tree
Changes from all 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
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,19 @@ bool SubscriptionEventSink::Subscription::Enqueue(const proto::optimizer::TraceQ
return true;
}

bool SubscriptionEventSink::Subscription::EnqueueBatch(
const std::vector<proto::optimizer::TraceQueryRequest> &events) {
{
std::lock_guard<std::mutex> lock(mutex_);
if (closed_ || events.size() > queue_size_ - queue_.size()) {
return false;
}
queue_.insert(queue_.end(), events.begin(), events.end());
}
cv_.notify_all();
return true;
}

void SubscriptionEventSink::Subscription::Close() {
{
std::lock_guard<std::mutex> lock(mutex_);
Expand Down Expand Up @@ -107,6 +120,7 @@ void SubscriptionEventSink::DisableSubscriptions() {
}

bool SubscriptionEventSink::Send(const proto::optimizer::TraceQueryRequest &event) {
std::lock_guard<std::mutex> send_lock(send_mutex_);
std::vector<std::shared_ptr<Subscription>> subscriptions;
{
std::lock_guard<std::mutex> lock(subscriptions_mutex_);
Expand All @@ -132,6 +146,37 @@ bool SubscriptionEventSink::Send(const proto::optimizer::TraceQueryRequest &even
return delivered;
}

bool SubscriptionEventSink::SendBatch(const std::vector<proto::optimizer::TraceQueryRequest> &events) {
if (events.empty()) {
return true;
}

std::lock_guard<std::mutex> send_lock(send_mutex_);
std::vector<std::shared_ptr<Subscription>> subscriptions;
{
std::lock_guard<std::mutex> lock(subscriptions_mutex_);
if (stopped_ || !accepting_subscriptions_) {
dropped_.fetch_add(events.size());
return false;
}
subscriptions = subscriptions_;
}
if (subscriptions.empty()) {
dropped_.fetch_add(events.size());
return false;
}

bool delivered = false;
for (const auto &subscription : subscriptions) {
if (subscription->EnqueueBatch(events)) {
delivered = true;
} else {
dropped_.fetch_add(events.size());
}
}
return delivered;
}

void SubscriptionEventSink::Stop() {
if (stopped_.exchange(true)) {
return;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ class SubscriptionEventSink final : public EventSink {

Subscription(std::string consumer_id, std::size_t queue_size);
bool Enqueue(const proto::optimizer::TraceQueryRequest &event);
bool EnqueueBatch(const std::vector<proto::optimizer::TraceQueryRequest> &events);
void Close();

const std::string consumer_id_;
Expand All @@ -61,6 +62,7 @@ class SubscriptionEventSink final : public EventSink {
bool accepting_subscriptions() const { return accepting_subscriptions_.load(); }

bool Send(const proto::optimizer::TraceQueryRequest &event) override;
bool SendBatch(const std::vector<proto::optimizer::TraceQueryRequest> &events);
void Stop() override;
std::size_t DroppedCount() const;

Expand All @@ -69,6 +71,7 @@ class SubscriptionEventSink final : public EventSink {

private:
OptimizerEventPublisherConfig config_;
std::mutex send_mutex_;
mutable std::mutex subscriptions_mutex_;
std::vector<std::shared_ptr<Subscription>> subscriptions_;
std::atomic<bool> accepting_subscriptions_{true};
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
24 changes: 24 additions & 0 deletions kv_cache_manager/event/test/subscription_event_sink_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
#include <gtest/gtest.h>
#include <memory>
#include <string>
#include <vector>

#include "kv_cache_manager/common/unittest.h"
#include "kv_cache_manager/event/optimizer_stream/subscription_event_sink.h"
Expand Down Expand Up @@ -71,6 +72,29 @@ TEST_F(SubscriptionEventSinkTest, TestDropsWithoutSubscriberAndAtQueueLimit) {
EXPECT_EQ("queued", event.trace_id());
}

TEST_F(SubscriptionEventSinkTest, TestBatchIsAllOrNothingPerSubscriber) {
SubscriptionEventSink sink(MakeSinkConfig(1, 2));
auto subscription = sink.Subscribe("optimizer");
ASSERT_NE(nullptr, subscription);

const std::vector<proto::optimizer::TraceQueryRequest> batch = {
MakeRequest("one"), MakeRequest("two")};
ASSERT_TRUE(sink.SendBatch(batch));

proto::optimizer::TraceQueryRequest event;
ASSERT_EQ(SubscriptionEventSink::Subscription::WaitResult::kEvent,
subscription->WaitNext(&event, std::chrono::milliseconds(10)));
EXPECT_EQ("one", event.trace_id());
ASSERT_EQ(SubscriptionEventSink::Subscription::WaitResult::kEvent,
subscription->WaitNext(&event, std::chrono::milliseconds(10)));
EXPECT_EQ("two", event.trace_id());

SubscriptionEventSink full_sink(MakeSinkConfig(1, 1));
ASSERT_NE(nullptr, full_sink.Subscribe("slow"));
EXPECT_FALSE(full_sink.SendBatch(batch));
EXPECT_EQ(2u, full_sink.DroppedCount());
}

TEST_F(SubscriptionEventSinkTest, TestRejectsExcessSubscribers) {
SubscriptionEventSink sink(MakeSinkConfig(1, 1));
ASSERT_NE(nullptr, sink.Subscribe("first"));
Expand Down
74 changes: 34 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,30 @@ 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",
":optimizer_call_guard",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_collector",
"//kv_cache_manager/optimizer/metrics:optimizer_metrics_reporter",
":optimizer_service_impl",
"//kv_cache_manager/common",
"//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 +179,15 @@ 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_kmonitor_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_group`、`instance_id` 和 `target_hit_rate_percent`,单位 byte)。`target_hit_rate_percent` 是相对于本上报窗口理论最大可命中量的比例,不是绝对请求命中率:目标命中量等于窗口理论无限容量最大可命中 block 数乘以该比例,`mrc` 则表示保留这些目标命中所需的最小 LRU 容量。当前固定输出 60%、80%、90%、95%、99%、99.5% 六个相对目标。例如理论最大命中率为 68.6% 时,95% 相对目标对应约 65.17% 的绝对命中率,而不是 95%。每次上报会原子取走并清空仅供 MRC 使用的 hit curve,不影响查询数、命中率等累计指标。该值直接聚合 LiteHit 产生的容量无关 hit curve,不依赖预先配置的离散容量点;周期内尚无理论可命中 block 时值为 0。


### Eviction Policies

Expand Down
Loading
Loading