代码:src/discovery/memory_registry.{h,cc} (180 行) + src/discovery/etcd_registry.{h,cc} (260 行)
场景:rpc_client 调 EchoService.Echo,注册中心返回所有 EchoService 实例的 host:port
RPC 调用第一步:“EchoService 的实例在哪些机器上?”——这就是服务发现。三种主流实现:
- 直连——客户端 hardcode host:port,简单但无扩展性
- 内存注册中心——单进程内 map<服务名, 实例列表>,适合开发/单测
- 独立服务——etcd / consul / nacos,跨进程共享
我两个都实现了。这篇文记录为什么需要两个、怎么设计、怎么从内存版升级到 etcd。
一、问题建模:三件事
不管用什么实现,服务发现都要做三件事:
1.1 Register:服务上线
1
| void register(string service, string host, uint16_t port, map<string,string> meta);
|
1.2 Heartbeat:服务保活
1
| void heartbeat(string service, string instance_id, int ttl_sec = 15);
|
1.3 Discover:客户端查询
1
2
| vector<Instance> discover(string service);
// 返回所有健康实例
|
关键点:Heartbeat + TTL——注册后 15s 没续约,自动剔除。
二、内存版:std::unordered_map + 心跳线程
2.1 数据结构
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
| class MemoryRegistry {
struct Instance {
string instance_id;
string host;
uint16_t port;
map<string,string> meta;
chrono::steady_clock::time_point last_heartbeat;
};
struct ServiceEntry {
unordered_map<string, Instance> instances; // instance_id → info
mutex mu;
};
unordered_map<string, ServiceEntry> services_; // service_name → instances
mutex meta_mu_; // 保护 services_ 这个 map
};
|
2.2 三个 API
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
| // 1. Register
void register(const string& service, const string& host, uint16_t port) {
string id = gen_uuid();
lock_guard<mutex> lk(meta_mu_);
auto& entry = services_[service];
lock_guard<mutex> lk2(entry.mu);
entry.instances[id] = {id, host, port, {}, Clock::now()};
}
// 2. Heartbeat
void heartbeat(const string& service, const string& instance_id) {
lock_guard<mutex> lk(meta_mu_);
auto it = services_.find(service);
if (it == services_.end()) return;
lock_guard<mutex> lk2(it->second.mu);
auto iit = it->second.instances.find(instance_id);
if (iit != it->second.instances.end()) {
iit->second.last_heartbeat = Clock::now();
}
}
// 3. Discover
vector<Instance> discover(const string& service) {
vector<Instance> result;
lock_guard<mutex> lk(meta_mu_);
auto it = services_.find(service);
if (it == services_.end()) return result;
lock_guard<mutex> lk2(it->second.mu);
for (auto& [id, inst] : it->second.instances) {
if (is_alive(inst)) result.push_back(inst);
}
return result;
}
|
2.3 心跳扫描线程
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
| void MemoryRegistry::start_health_checker() {
health_thread_ = thread([this] {
while (!stop_) {
this_thread::sleep_for(5s);
evict_expired();
}
});
}
void MemoryRegistry::evict_expired() {
auto now = Clock::now();
lock_guard<mutex> lk(meta_mu_);
for (auto& [service, entry] : services_) {
lock_guard<mutex> lk2(entry.mu);
for (auto it = entry.instances.begin(); it != entry.instances.end();) {
if (now - it->second.last_heartbeat > 15s) {
log("evict %s/%s (no heartbeat)", service, it->first);
it = entry.instances.erase(it);
} else {
++it;
}
}
}
}
|
5s 扫一次,15s 没心跳就剔除——简单粗暴。
2.4 内存版的局限
致命:只能单进程。两个 RPC server 在不同机器,各自的 MemoryRegistry 互相不知道。
所以内存版只用在:
- 单测
- 单进程 demo
- 嵌入式场景(全部代码跑在一个 binary 里)
生产必须 etcd/consul。
三、etcd 版:Lease + Watch + gRPC
3.1 etcd 的三个核心概念
| 概念 | 作用 | 类比 |
|---|
| Key-Value | 数据存储 | Redis hash |
| Lease | 租约(TTL),过期自动删 key | Redis EXPIRE |
| Watch | 订阅 key 变化 | Redis keyspace notification |
我用 etcd 存的 key 格式:
1
2
| /services/<service_name>/<instance_id> = "host:port"
+ JSON meta
|
3.2 三个 API 的 etcd 实现
Register:用 Lease
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
| // etcd v3 gRPC API: LeaseGrant + Put
void EtcdRegistry::register(const string& service, const string& host, uint16_t port) {
// 1. 申请 lease (15s TTL)
auto lease_resp = lease_stub_->LeaseGrant(LeaseGrantRequest{15});
string lease_id = to_string(lease_resp.lease_id());
// 2. put key 带 lease
string key = "/services/" + service + "/" + gen_uuid();
string value = host + ":" + to_string(port);
put_stub_->Put(PutRequest{key, value, lease_id});
// 3. 后台线程续约 (每 5s 续一次)
thread([this, lease_id] {
while (!stop_) {
this_thread::sleep_for(5s);
lease_stub_->LeaseKeepAlive(LeaseKeepAliveRequest{lease_id});
}
}).detach();
}
|
Key 关键:LeaseKeepAlive 续约——etcd 内部 keepalive,client 只要发"续"就续。
Heartbeat:就是 LeaseKeepAlive
etcd 心跳 = lease 续约,不需要单独 API。
Discover:用 Range + Watch
1
2
3
4
5
6
7
8
9
10
| // 1. 首次发现:Range 查
auto resp = kv_stub_->Range(RangeRequest{
key: "/services/" + service + "/",
range_end: "/services/" + service + "/\xff" // 前缀匹配
});
vector<Instance> result;
for (auto& kv : resp.kvs()) {
result.push_back(parse(kv.value()));
}
return result;
|
/services/EchoService/\xff——\xff 是比任何字符都大的字节,Range key 在 [start, end) 之间,前缀匹配的写法。
1
2
3
4
5
6
7
8
9
10
11
| // 2. 订阅变化:Watch
void EtcdRegistry::watch(const string& service, function<void(vector<Instance>)> cb) {
watch_stub_->Watch(WatchRequest{
key: "/services/" + service + "/"
}, [cb](auto resp) {
if (resp.events_size() > 0) {
// 重新拉一次 Range
cb(discover(service));
}
});
}
|
客户端不需要轮询——etcd 主动推变化,延迟 < 100ms。
3.3 etcd 版的优势
| 维度 | 内存版 | etcd 版 |
|---|
| 跨进程 | ❌ | ✅ |
| 持久化 | ❌(进程挂了注册就没了) | ✅ |
| Watch 推送 | ❌(只能轮询) | ✅ |
| 强一致 | ❌ | ✅(Raft) |
| 复杂度 | 180 行 | 260 行 |
| 依赖 | 无 | etcd cluster (3 节点起) |
3.4 etcd 版的坑
坑 1:Lease 续约失败
LeaseKeepAlive 失败时,lease 还在,但续约已断——15s 后 etcd 自动删 key。
解决:续约失败要重新 register——我加了"如果连续 2 次 keepalive 失败,重置 lease":
1
2
3
4
5
6
7
8
9
10
11
12
13
14
| int retry = 0;
while (!stop_) {
this_thread::sleep_for(5s);
try {
lease_stub_->LeaseKeepAlive(lease_id);
retry = 0;
} catch (...) {
if (++retry > 2) {
log("lease %s lost, re-registering", lease_id);
register(service, host, port); // 重置
break;
}
}
}
|
坑 2:Watch 流断了要重连
gRPC stream 一断,不会自动重连。我用 grpc::ClientContext 每次重连:
1
2
3
4
5
6
7
8
| while (!stop_) {
ClientContext ctx;
WatchRequest req{key};
WatchResponse resp;
watch_stub_->Watch(&ctx, req, ...); // 阻塞到断开
log("watch stream broken, reconnecting in 1s");
this_thread::sleep_for(1s);
}
|
坑 3:Range 大量 key 性能
1000 个实例时,Range 返回 1000 条 KV——单次 RPC 几 MB。
解决:用 Limit + 分页:
1
2
3
4
5
| RangeRequest req;
req.set_key(prefix);
req.set_range_end(prefix_end);
req.set_limit(100); // 每页 100
req.set_sort_target(REQUESTHASH); // 稳定排序
|
四、客户端集成:Discovery 抽象
RpcClient 不应该绑死"用 etcd 还是内存"——加抽象:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
| class ServiceDiscovery {
public:
virtual vector<Instance> discover(const string& service) = 0;
virtual void watch(const string& service,
function<void(vector<Instance>)> cb) = 0;
};
class RpcClient {
shared_ptr<ServiceDiscovery> discovery_;
// 内部:30s 缓存,避免每次 RPC 都查注册中心
vector<Instance> get_instances(const string& service) {
auto& cached = cache_[service];
if (cached.instances.empty() ||
Clock::now() - cached.fetched_at > 30s) {
cached.instances = discovery_->discover(service);
cached.fetched_at = Clock::now();
}
return cached.instances;
}
};
|
30s 缓存——平衡"实时性"和"避免对 etcd 压力"。
首次 watch:
1
2
3
4
5
| // RpcClient 构造时
discovery_->watch("EchoService", [this](auto new_list) {
cache_["EchoService"].instances = new_list;
cache_["EchoService"].fetched_at = Clock::now();
});
|
Push vs Pull 协同:
- 首次:
discover() pull 一次 - 之后:etcd push 触发 watch 回调,直接更新缓存
- 30s 没 push(etcd 故障兜底):下次
discover() pull 一次
五、生产实践:3 个选择
5.1 单测 / Demo:用 MemoryRegistry
1
2
3
4
5
6
7
| TEST_CASE("client picks healthy instance") {
MemoryRegistry reg;
reg.register("EchoService", "10.0.0.1", 9000);
reg.register("EchoService", "10.0.0.2", 9000);
RpcClient client(make_shared<MemoryRegistryAdapter>(reg));
// ...
}
|
5.2 中小规模(< 100 实例):etcd
3 节点 etcd cluster,足够。注意 etcd 的 fsync 性能——机械盘会卡,用 SSD。
5.3 大规模(> 1000 实例):nacos / Polaris / 自研
etcd 的 KV 数量超过 10w 时,Range 性能下降。大集群换 nacos(AP) 或自研 AP 方案(牺牲强一致换可用性)。
六、最后一个反直觉的点:服务发现不是免费的
服务发现的代价:
- 每实例 5s 一次心跳——1000 实例 = 200 qps 持续打 etcd
- Watch stream 维护——1000 服务订阅 = 1000 个长连接
- 缓存一致性——etcd 推过来 100ms 延迟,客户端在期间可能拿到旧列表
所以:不是所有服务都要进注册中心。
- 静态配置的服务(数据库、Redis)直接读 config file
- 短生命周期任务(批处理)直接连
- 真正需要"动态扩缩容"的核心服务(用户服务、订单服务)才进 etcd
别为了"用上 etcd"而把 5 个服务的注册都堆进去。
相关阅读: