0%

go-zero 源码分析 09:服务发现与负载均衡

在上一篇文章的最后,我们追踪了一次完整的 zRPC 调用——客户端拦截器从 Tracing 走到 Timeout,然后在 invoker 中将请求发给 gRPC 底层。但这里跳过一个至关重要的问题:invoker 发请求时,到底发给了谁?

这个问题看似简单。直连模式最简单——配置里写死了 127.0.0.1:8080,直接连就行。但生产环境不可能用直连。服务实例分布在多台机器、多个容器中,IP 地址随 Pod 调度而漂移,端口也可能因为 Sidecar 注入而变化。客户端必须自动发现正确的目标地址列表,并在多个地址之间做出选择。

这就是本文要解决的三个核心问题:

  1. 服务发现:当服务端启动后,客户端如何知道它的存在?当服务端下线后,客户端如何感知并剔除它的地址?
  2. 地址更新:地址列表的动态变更是如何从 etcd 的 watch channel 一路传递到 gRPC 的 connection 层面的?
  3. 负载均衡:面对 N 个可用节点,客户端选哪个?怎么判断一个节点"好不好"?算法如何平衡负载分配与单点过热?

三者的关系是一条完整的数据流:

1
2
3
服务端启动 → Publisher 注册到 etcd → Subscriber Watch 到变化
→ Resolver 更新 gRPC 地址列表 → Balancer 为每个地址建立 SubConn
→ 每次 RPC 调用时 P2C Picker 选择一个 SubConn → 请求发出

下面我们就沿着这条链路,从服务注册出发,经过 watch 状态机和 resolver 转换,最终到达 P2C 算法的核心。

服务注册:Publisher 如何把"我在线"告诉 etcd

问题:服务端重启后 IP 变了怎么办

一个朴素的做法是运维人员手工更新配置中的服务地址列表然后重启客户端。但在微服务架构中,服务实例的数量和位置是随时间动态变化的——扩容、缩容、滚动更新、故障恢复,每一次变化都意味着地址的改变。你不可能靠手工跟上这个节奏。

go-zero 的答案清晰而直接:服务端自己告诉注册中心"我在哪",客户端从注册中心订阅"有哪些"。在 go-zero 中,这个注册中心就是 etcd——一个强一致性的分布式键值存储。服务端启动时,将自己的地址作为 key-value 写入 etcd 的特定前缀下,并以租约(Lease)保持心跳;租约过期,etcd 自动删除这条记录,客户端看到的就是最新的在线实例列表。

从 KeepAliveServer 说起

在之前的服务端启动流程中,我们讲到当 RpcServerConf 配置了 etcd 时,NewServer 会创建 RpcPubServer 而非普通的 RpcServerRpcPubServerStart 时额外执行了一步 registerEtcd

1
2
3
4
5
6
7
// zrpc/internal/rpcpubserver.go
func (s keepAliveServer) Start(fn RegisterFn) error {
if err := s.registerEtcd(); err != nil {
return err
}
return s.Server.Start(fn)
}

registerEtcd 的核心是问 discov.Publisher 注册当前服务实例:

1
2
3
4
5
6
7
// zrpc/internal/rpcpubserver.go (registerEtcd 内部)
pub := discov.NewPublisher(
s.config.Etcd.Hosts,
s.config.Etcd.Key,
listenOn, // 经过 figureOutListenOn 处理后的真实可路由地址
)
pub.KeepAlive()

这完成了两条信息的关联:在 etcd 的 s.config.Etcd.Key 前缀下,注册了一个 key,它的 value 是服务进程的监听地址。 etcd 中实际存储的结构类似:

1
2
Key:   /service/greet.rpc/7587844881234567890
Value: 192.168.1.10:8080

其中 key 的最后一串数字是租约 ID(如果没有指定自定义 ID),value 是经过 figureOutListenOn 转换后的真实地址——不是 0.0.0.0:8080,而是从环境变量 POD_IP 或网卡检测得到的真实 IP。如果 Key 对应的 Value 被其他实例先前注册过,Publisher 的 Exclusive 模式会把旧 Key 移除,保证每个地址只注册一次。

KeepAlive 的三个并发循环

pub.KeepAlive() 启动了一个异步的 goroutine(通过 threading.GoSafe),在这个 goroutine 中同时维护三个 channel 的监听:

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
34
35
36
37
38
39
40
41
42
// core/discov/publisher.go
func (p *Publisher) keepAliveAsync(cli internal.EtcdClient) error {
ch, err := cli.KeepAlive(cli.Ctx(), p.lease)
// ...
threading.GoSafe(func() {
wch := cli.Watch(cli.Ctx(), p.fullKey, clientv3.WithFilterPut())

for {
select {
case _, ok := <-ch:
// 租约心跳 channel 关闭 → 重新注册
if !ok {
p.revoke(cli)
if err := p.doKeepAlive(); err != nil { ... }
return
}
case c := <-wch:
// etcd 上自己的 key 被删除 → 重新 put
for _, evt := range c.Events {
if evt.Type == clientv3.EventTypeDelete {
cli.Put(cli.Ctx(), p.fullKey, p.value, clientv3.WithLease(p.lease))
}
}
case <-p.pauseChan:
// 暂停续约
p.revoke(cli)
select {
case <-p.resumeChan:
p.doKeepAlive()
return
case <-p.quit.Done():
return
}
case <-p.quit.Done():
// 优雅退出
p.revoke(cli)
return
}
}
})
return nil
}

三个并发循环各司其职:

  • 租约心跳 channel:etcd 为每个 key 分配了一个带 TTL 的租约(默认 10 秒),cli.KeepAlive 返回一个只读 channel,etcd 客户端定期向这个 channel 发送心跳续约事件。如果这个 channel 被关闭(可能因为网络分区或 etcd 服务不可达),Publisher 立即撤销旧租约并执行 doKeepAlive——每秒重试一次,直到重新注册成功。
  • Watch channel:即使心跳正常,其他外部因素(如管理员手工删除了 key)也可能导致 key 消失。Publisher 用 clientv3.WithFilterPut() Watch 自己的 key,过滤掉 PUT 事件(自己 put 的,不需要响应),只监听 DELETE——一旦发现被删了,立即重新 put。
  • Pause/Resume/Quit:提供了运行时的生命周期控制。Pause 撤销租约让 key 从服务发现中消失(常用于流量摘除的场景),Resume 重新注册继续接收流量。

这种"三通道监听"设计的根本原因是:每一种失败模式都需要独立的感知和恢复策略。 心跳断了可能是网络闪断(重试即可),key 被删了可能是外部修改(需要立即重建),暂停和退出是主动行为(需要精确的执行路径)。把所有处理塞进一个 recover 包裹的循环里固然能兜底,但丢失了每种场景的准确语义。

重试的冷静期:CoolDown 机制

注意到租约断开后的 doKeepAlive 并不是疯狂无限重试——它每秒执行一次。而在 clusterwatch 重试循环中,还有一个更精细的机制:

1
2
3
4
5
6
7
8
9
10
11
// core/discov/internal/registry.go
func (c *cluster) watch(cli EtcdClient, key watchKey, rev int64) {
for {
err := c.watchStream(cli, key, rev)
if err == nil {
return
}
// ...
time.Sleep(coolDownUnstable.AroundDuration(coolDownInterval))
}
}

coolDownUnstable 是一个 mathx.Unstable 实例,偏差为 0.05(5%),基础间隔为 1 秒。所以实际重试间隔在 [0.95s, 1.05s] 之间随机。为什么需要一个带随机因子的冷却间隔?

设想这样一个场景:100 个服务实例同时从 etcd 断开,它们全部在同一毫秒重连——这个"惊群效应"会让 etcd 瞬间压力剧增,可能触发二次断开,形成雪崩。引入 5% 的随机抖动后,每个实例的重连时间会均匀分散在 100ms 的窗口内,etcd 的负载被摊平了。这是一个几十行代码的小机制,但反映了 go-zero 对分布式系统级联风险的系统性思考。

服务发现:Subscriber 如何感知实例变化

问题:如何在全量 return 和增量 push 之间找到平衡

服务端注册完成后,客户端需要感知到"有哪些服务端在线"。常见的做法有两类:

  • 轮询:客户端定时拉取全量服务列表。简单但浪费——无论列表变没变都要查,实例数多时尤其低效。
  • 长轮询/Watch:客户端订阅 key 的变化,只在变更时才收到通知。高效但有状态——需要维护连接、处理断连重连。

go-zero 用的是 etcd 原生的 Watch 机制:客户端发起一个 Get(全量同步初始值),再启动一个 Watch(增量接收后续变更),两者之间通过 revision 机制无缝衔接。

Subscriber 的创建与容器的绑定

Subscriber 的使用模式很简单:指定 etcd 地址和要订阅的 key 前缀,然后通过 Values() 获取当前在线地址列表:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// core/discov/subscriber.go
func NewSubscriber(endpoints []string, key string, opts ...SubOption) (*Subscriber, error) {
sub := &Subscriber{
endpoints: endpoints,
key: key,
}
for _, opt := range opts {
opt(sub)
}
if sub.items == nil {
sub.items = newContainer(sub.exclusive)
}

if err := internal.GetRegistry().Monitor(endpoints, key, sub.exactMatch, sub.items); err != nil {
return nil, err
}

return sub, nil
}

internal.GetRegistry() 返回的是一个全局单例Registry——同一进程内所有 Subscriber 和 Publisher 共享。Monitorsub.items(一个 Container 实例)注册为指定 key 的监听器。这里的核心抽象是:

1
2
3
4
5
// core/discov/internal/updatelistener.go
type UpdateListener interface {
OnAdd(kv KV)
OnDelete(kv KV)
}

每当 etcd 上该 key 前缀下的记录有新增或删除,OnAddOnDelete 就会被调用。而 Container(即 container 实现)正是这个接口的实现者——它维护着一份本地的 values map,反映当前在线实例的完整集合。

Container:本地缓存的增量更新

container 保存了两张内存表:

1
2
3
4
5
6
7
8
9
10
// core/discov/subscriber.go
type container struct {
exclusive bool
values map[string][]string // value → [key1, key2]
mapping map[string]string // key → value
snapshot atomic.Value // 快照:[]string
dirty *syncx.AtomicBool
listeners []func()
lock sync.Mutex
}

values 的 key 是服务地址(value),value 是 etcd key 列表。一个地址可能对应多个 etcd key(不同租约 ID),exclusive 模式下新 key 注册时会清除同 value 的旧 key。

snapshotdirty 实现了一个惰性快照机制。GetValues()dirty 为 false 时直接返回上次生成的快照(零开销的原子读),只在 dirty 为 true 时才加锁重建。这个优化看似微小,但在高并发的 Values() 调用(每次 pick 前都要获取地址列表)中,从加锁遍历 map 变成了原子值读取,性能差异显著。

notifyChange 更为关键——每次 OnAddOnDelete 完成后,它遍历所有注册的 listener 并调用:

1
2
3
4
5
6
7
8
9
func (c *container) notifyChange() {
c.lock.Lock()
listeners := append(([]func())(nil), c.listeners...)
c.lock.Unlock()

for _, listener := range listeners {
listener()
}
}

这个 listener 是怎么注册的?答案就在 resolver 的 Build 方法中。

Resolver:从 etcd 变化到 gRPC 地址的桥梁

container 帮我们解决了"本地维护实例列表"的问题,但 gRPC 需要的是 resolver.Address 列表。gRPC 的 resolver 接口定义了两个核心职责:解析 target 字符串,以及通过 resolver.ClientConn.UpdateState 告知 gRPC 地址变化。go-zero 要做的是:在 Subscriber 的 listener 和 cc.UpdateState 之间架一座桥。

四种 Resolver Scheme

go-zero 在 init() 中注册了四种 resolver scheme:

  • direct:直连模式。target 如 direct:///127.0.0.1:8080,127.0.0.1:8081。地址在 Build 时一次性解析完毕,返回一个 nopResolver——不需要 watch,不需要更新。
  • discov:服务发现模式(etcd watch)。target 如 discov:///etcd-host:2379?key=greet.rpcBuild 时创建 Subscriber,注册 listener,地址变化时调用 cc.UpdateState
  • etcd:与 discov 功能完全相同,只是 scheme 名不同。etcdBuilder 内嵌了 discovBuilder,仅覆盖了 Scheme() 方法返回 "etcd"。这是为了兼容历史配置。
  • k8s:Kubernetes 模式。target 如 k8s:///namespace/servicename:port。不通过 etcd,而是使用 Kubernetes 原生的 Informer 机制监听 EndpointSlice 资源变化。

关键的桥接:listener 和 cc.UpdateState

discovBuilder.Build 只有三十几行,但它是整条链路中最关键的桥接点:

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
// zrpc/resolver/internal/discovbuilder.go
func (b *discovBuilder) Build(target resolver.Target, cc resolver.ClientConn, _ resolver.BuildOptions) (
resolver.Resolver, error) {
hosts := strings.FieldsFunc(targets.GetHosts(target), func(r rune) bool {
return r == EndpointSepChar
})
sub, err := discov.NewSubscriber(hosts, targets.GetKey(target))
if err != nil {
return nil, err
}

update := func() {
vals := subset(sub.Values(), subsetSize)
addrs := make([]resolver.Address, 0, len(vals))
for _, val := range vals {
addrs = append(addrs, resolver.Address{
Addr: val,
})
}
if err := cc.UpdateState(resolver.State{
Addresses: addrs,
}); err != nil {
logx.Error(err)
}
}
sub.AddListener(update)
update()

return &discovResolver{
cc: cc,
sub: sub,
}, nil
}

让我们仔细追踪这个函数的行为:

  1. 解析 targettargets.GetHoststargets.GetKey 从 gRPC target URI 中提取 etcd 主机列表和订阅的 key 前缀。新版格式是 discov:///h1:port,h2:port?key=greet.rpc——hosts 放在 URI path 中以避免逗号分隔的多主机地址违反 RFC 3986;旧版格式 etcd://h:port/key 仍兼容。

  2. 创建 Subscriber:传入 etcd 地址列表和要订阅的 key 前缀。Subscriber 内部立即执行一次全量 Get 后启动 Watch

  3. 注册 listenerupdate 闭包函数做了三件事——从 Container 获取当前所有在线地址(sub.Values()),通过 subset 取最多 32 个(防止一个服务有数百个实例时 gRPC 管理的连接数爆炸),构造 resolver.Address 列表,调用 cc.UpdateState 提交给 gRPC。

  4. 立即调用一次 update:在注册完 listener 后立即调用,将初始地址列表推给 gRPC。这一步保证了服务刚启动时客户端就能拿到地址,而不是等到下一次 watch 事件。

从此刻起,数据流变成:

1
2
3
etcd watch 事件 → cluster.handleWatchEvents → UpdateListener.OnAdd/OnDelete
→ container addKv/removeKey → notifyChange → listener (update 闭包)
→ subset → cc.UpdateState → gRPC balancer 收到新地址

每一次 etcd 上的变化,都会自动推送到 gRPC 的连接管理层面。 新增实例会触发新 SubConn 的创建,下线实例会触发旧 SubConn 的关闭和回收——整个过程对业务代码完全透明。

连接复用的隐形架构

在 Subscriber 和 Publisher 的背后,有一个容易忽视但至关重要的基础设施:Registry 中 etcd 客户端连接的复用。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// core/discov/internal/registry.go
func (r *Registry) GetConn(endpoints []string) (EtcdClient, error) {
c, _ := r.getOrCreateCluster(endpoints)
return c.getClient()
}

func (c *cluster) getClient() (EtcdClient, error) {
val, err := connManager.GetResource(c.key, func() (io.Closer, error) {
return c.newClient()
})
if err != nil {
return nil, err
}
return val.(EtcdClient), nil
}

connManager 是一个 syncx.ResourceManager——它将 etcd 端点列表排序后拼接为唯一 key,相同端点列表只创建一个 etcd 客户端连接。这意味着:

  • 同一个进程内对同一个 etcd 集群的多个 Subscriber 和 Publisher 共享一个连接。不会因为你订阅了 10 个不同的 key 前缀就创建 10 个 gRPC 连接到 etcd。
  • 连接创建由 ResourceManager 保证并发安全——同时有多个 goroutine 请求连接时,只有一个真正创建,其他等待并复用结果(这实际上是 ResourceManager 内部使用了 SingleFlight 机制)。

etcd 连接的自动恢复

共享连接带来了连接数优化的好处,但也引入了一个问题:如果底层连接断开后重建,所有 Watch 全部失效,谁来重启它们?

答案在 StateWatcher

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
// core/discov/internal/statewatcher.go
func (sw *stateWatcher) watch(conn etcdConn) {
sw.currentState = conn.GetState()
for {
if conn.WaitForStateChange(context.Background(), sw.currentState) {
sw.updateState(conn)
}
}
}

func (sw *stateWatcher) updateState(conn etcdConn) {
sw.currentState = conn.GetState()
switch sw.currentState {
case connectivity.TransientFailure, connectivity.Shutdown:
sw.disconnected = true
case connectivity.Ready:
if sw.disconnected {
sw.disconnected = false
sw.notifyListeners()
}
}
}

StateWatcher 监听了 etcd gRPC 连接的状态变化。当连接从 TransientFailureShutdown 恢复为 Ready 时,它会通知 cluster.reload——后者会取消所有旧的 Watch context,等待旧的 watch goroutine 退出,然后启动全新的 Get + Watch 循环。

1
2
3
4
5
6
7
8
// core/discov/internal/registry.go
func (c *cluster) watchConnState(cli EtcdClient) {
watcher := newStateWatcher()
watcher.addListener(func() {
go c.reload(cli)
})
watcher.watch(cli.ActiveConnection())
}

这个设计保证了:即便 etcd 集群整体故障后恢复,客户端的服务发现能力也能自动恢复。 不需要人工重启客户端进程,不需要额外的健康检查定时器——连接恢复的信号会沿着 StateWatcher → cluster.reload → 重新 Get + Watch → UpdateListener → cc.UpdateState 一路传导到 gRPC 的负载均衡层。

Revision 压缩的处理

etcd 为了控制存储大小,会定期压缩历史 revision。如果客户端的 Watch 起始 revision 已经被压缩,Watch 会返回 rpctypes.ErrCompacted 错误。

1
2
3
4
5
6
7
8
9
10
11
12
13
func (c *cluster) watch(cli EtcdClient, key watchKey, rev int64) {
for {
err := c.watchStream(cli, key, rev)
if err == nil {
return
}
if rev != 0 && errors.Is(err, rpctypes.ErrCompacted) {
logc.Errorf(cli.Ctx(), "etcd watch stream has been compacted, try to reload, rev %d", rev)
rev = c.load(cli, key) // 重新全量 load,获取最新 revision
}
time.Sleep(coolDownUnstable.AroundDuration(coolDownInterval))
}
}

当 revision 被压缩时,框架不会报错放弃,而是重新执行 load(全量 Get)获取最新 revision,以此为基础重新启动 Watch。这保证了在 etcd 运维期间的短暂不可用不会导致客户端的地址列表永久停滞。

全局视图:Registry 与 Cluster 的生命周期

在进入负载均衡之前,先纵览一下服务发现的整体架构。Registry 内部管理的是一个 map[string]*cluster——同一组 etcd 端点对应一个 cluster

1
2
3
4
5
6
7
8
type cluster struct {
endpoints []string
key string
watchers map[watchKey]*watchValue
watchGroup *threading.RoutineGroup
done chan lang.PlaceholderType
lock sync.RWMutex
}

watchKey{key, exactMatch} 组成,watchValue 包含该 key 上所有 listener 和当前的全量 values 缓存。不同的 watchKey 共享同一个 etcd 连接,但各自独立维护自己的 Watch 协程和 values 缓存。 例如,一个客户端同时订阅了 /service/order.rpc/service/user.rpc,它们注册到同一个 cluster 下,共享一个 etcd 客户端连接,但各自有独立的 watch goroutine。

watchGroup 是一个 RoutineGroup,管理了所有 watch goroutine 的生命周期。cluster.reload 在执行时,先关闭 done channel(退出所有 watch goroutine),调用 watchGroup.Wait() 等待全部退出,然后重置 donewatchGroup,重新启动所有 key 的 watch。这个优雅的重启流程避免了 goroutine 泄漏和并发冲突。

subscriber.Close() 对应的 Registry.Unmonitor 流程也是一样自洽的:从 listener 列表中移除当前 listener,如果这是最后一个 listener,取消对应的 watch context 并从 c.watchers 中删除该 watchKey——没有 listener 还开着 watch 是浪费资源和带宽。

负载均衡:从地址列表到选择节点

问题:N 个节点,该选哪一个

Resolver 将地址列表交给了 gRPC,gRPC 的 balancer 会为每个地址建立一个 SubConn(子连接)。接下来,每次 RPC 调用时,Picker.Pick 需要从所有 Ready SubConn 中选择一个。

最简单的策略是轮询——每个请求选下一个节点,保证"所有节点被调用的次数相同"。但轮询忽略了一个关键事实:节点的处理能力并不相同。 一台刚恢复的节点可能还在预热(JIT、缓存填充),一台正在做 GC 的节点可能延迟很高,一台物理机上的实例和一台性能较差的虚拟机上的实例处理能力天然不同。

go-zero 的选择是 P2C(Power of Two Choices)+ EWMA(指数加权移动平均延迟)。这是一个在生产环境中经过充分验证的组合策略。

P2C 算法:两次随机而非一次选择

P2C 的全称是 “Power of Two Choices”——从列表中随机选择两个节点,取其中负载较轻的那个。这不是随机选一个,也不是遍历所有节点找出最优的(O(n)),而是用 O(1) 的随机和一次比较,理论上将最大负载降低到大约 log(log(n)) 级别。

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
34
35
36
37
38
// zrpc/internal/balancer/p2c/p2c.go
func (p *p2cPicker) Pick(_ balancer.PickInfo) (balancer.PickResult, error) {
p.lock.Lock()
defer p.lock.Unlock()

var chosen *subConn
switch len(p.conns) {
case 0:
return emptyPickResult, balancer.ErrNoSubConnAvailable
case 1:
chosen = p.choose(p.conns[0], nil)
case 2:
chosen = p.choose(p.conns[0], p.conns[1])
default:
var node1, node2 *subConn
for i := 0; i < pickTimes; i++ {
a := p.r.Intn(len(p.conns))
b := p.r.Intn(len(p.conns) - 1)
if b >= a {
b++
}
node1 = p.conns[a]
node2 = p.conns[b]
if node1.healthy() && node2.healthy() {
break
}
}
chosen = p.choose(node1, node2)
}

atomic.AddInt64(&chosen.inflight, 1)
atomic.AddInt64(&chosen.requests, 1)

return balancer.PickResult{
SubConn: chosen.conn,
Done: p.buildDoneFunc(chosen),
}, nil
}

pickTimes 等于 3——最多重试 3 次随机选出两个"健康"的节点。什么是健康?healthy() 的判断条件只有一行:

1
2
3
func (c *subConn) healthy() bool {
return atomic.LoadUint64(&c.success) > throttleSuccess
}

c.success 是一个每轮完成时通过 EWMA 更新的"成功率"指标,初始值为 1000(initSuccess),throttleSuccess 为 500(即初始成功率的 50%)。当节点的成功率低于一半时,它被认为"不健康",会被随机选择过程排除在外——但不是永远排除,稍后会讲到。

排他式随机选取

a := p.r.Intn(len(p.conns))b := p.r.Intn(len(p.conns) - 1) 然后 if b >= a { b++ },这是经典的"从 n 个元素中无放回选 2 个"的写法——a[0, n-1] 中选,b 从去掉 a 后的 n-1 个位置中选。短小的算法保证了 node1 和 node2 永远不会是同一个实例。

choose:负载比较与强制拾取

两个节点选出来了,怎么判断谁更"轻"?

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
func (p *p2cPicker) choose(c1, c2 *subConn) *subConn {
start := int64(timex.Now())
if c2 == nil {
atomic.StoreInt64(&c1.pick, start)
return c1
}

if c1.load() > c2.load() {
c1, c2 = c2, c1
}

// 强制拾取:如果 c2 被"冷落"太久,放弃 c1,选 c2
pick := atomic.LoadInt64(&c2.pick)
if start-pick > forcePick && atomic.CompareAndSwapInt64(&c2.pick, pick, start) {
return c2
}

atomic.StoreInt64(&c1.pick, start)
return c1
}

正常情况下,choose 选取 load() 较小的节点。load() 的计算是 P2C 算法的灵魂:

1
2
3
4
5
6
7
8
func (c *subConn) load() int64 {
lag := int64(math.Sqrt(float64(atomic.LoadUint64(&c.lag) + 1)))
load := lag * (atomic.LoadInt64(&c.inflight) + 1)
if load == 0 {
return penalty // penalty = math.MaxInt32
}
return load
}

load = sqrt(lag) × (inflight + 1),总结起来就是一句话:节点的"负载"由它的历史延迟和当前正在处理的请求数共同决定。

  • lag 取平方根是为了压缩延迟的数值范围——延迟从 1ms 变为 100ms 是非常大的恶化,从 100ms 变为 200ms 虽然也糟,但没有那么戏剧性地糟。平方根让 lag 对负载的影响递减而不是线性上升。
  • inflight + 1 保证了即使没有在途请求,负载也不为 0。
  • 如果 load 计算结果为 0(只有在 lag == 0 && inflight == -1 才可能),返回 math.MaxInt32——这个节点被视为"几乎不可能被选中",直到下次 lag 更新。

还有一个巧妙的机制:强制拾取。如果某个节点超过 1 秒(forcePick)没有被选中,choose 会忽略负载比较,直接选中它。为什么需要这个?

因为 EWMA 的延迟值是滞后的——如果某个节点的瞬时延迟很低(因为最近根本没有请求打到它),它的 lag 就会很低,load() 也很低,P2C 总是选中它,形成一个"富者越富"的正反馈循环。强制拾取打破了这种惯性——至少在每秒一次的频率上,每个节点都有机会被选到。这维持了所有 SubConn 的"热度",保证了所有的连接都保持了活跃的状态跟踪,也避免了延迟信号在长时间无请求后失去参考价值。

EWMA 延迟:为何指数加权

P2C 解决了"选谁"的问题。但如何衡量一个节点的延迟不是一次取样的暴政,而是对历史趋势的平滑反映?每次 RPC 调用完成后,Done 回调更新节点的 lag 和 success:

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
34
func (p *p2cPicker) buildDoneFunc(c *subConn) func(info balancer.DoneInfo) {
start := int64(timex.Now())
return func(info balancer.DoneInfo) {
atomic.AddInt64(&c.inflight, -1)
now := timex.Now()
last := atomic.SwapInt64(&c.last, int64(now))
td := int64(now) - last
if td < 0 {
td = 0
}

// EWMA 权重:时间间隔越大,历史数据的权重越低
w := math.Exp(float64(-td) / float64(decayTime))

lag := int64(now) - start
if lag < 0 {
lag = 0
}
olag := atomic.LoadUint64(&c.lag)
if olag == 0 {
w = 0
}

// 新 lag = 旧 lag × w + 当前 lag × (1-w)
atomic.StoreUint64(&c.lag, uint64(float64(olag)*w + float64(lag)*(1-w)))

success := initSuccess
if info.Err != nil && !codes.Acceptable(info.Err) {
success = 0
}
osucc := atomic.LoadUint64(&c.success)
atomic.StoreUint64(&c.success, uint64(float64(osucc)*w + float64(success)*(1-w)))
}
}

w 不是固定常数——它基于 距离上一次 Done 回调的时间间隔 td 动态计算:

  • 如果请求密集(td 很小),w 接近 1,历史数据保留更多权重,延迟估计更平滑。
  • 如果请求稀疏(td 很大),w 接近 0,历史数据几乎被丢弃,直接用当前延迟重新初始化——长时间无请求后,之前的延迟数据已经没有参考意义。

decayTime 固定为 10 秒——这是一个从 Twitter Finagle 继承的经验值。这个值的含义是:10 秒内的历史延迟具有实质参考意义,超过 10 秒则衰减到接近零。

这里有一个容易被忽略的细节:olag == 0w = 0。为什么?一个刚创建的 subConn,其 lag 初始值为 0。如果 w 不为 0,第一次 Done 回调后,新 lag 会是 0 * w + realLag * (1-w)——被历史值 0 拉低了。不合理。设置 w = 0 等价于直接使用第一次的实际延迟值作为起点,之后才开始平滑。

success 同理:codes.Acceptable 的判断(在第 08 篇中有详细分析)决定了这次调用是否算"成功"——DeadlineExceededInternalUnavailable 等错误算失败(success=0),InvalidArgumentNotFound 等错误算成功(success=1000)。通过 EWMA,success 形成了一个在 [0, 1000] 之间的连续值,低于 500 的节点被视为不健康。

统计日志:每分钟一张健康快照

还有一个容易被忽视但运维友好的功能——每分钟打印一次节点统计:

1
2
3
4
5
6
7
8
9
10
func (p *p2cPicker) logStats() {
stats := make([]string, 0, len(p.conns))
p.lock.Lock()
defer p.lock.Unlock()
for _, conn := range p.conns {
stats = append(stats, fmt.Sprintf("conn: %s, load: %d, reqs: %d",
conn.addr.Addr, conn.load(), atomic.SwapInt64(&conn.requests, 0)))
}
logx.Statf("p2c - %s", strings.Join(stats, "; "))
}

每分钟,你可以从日志中看到每个节点的当前负载值和过去一分钟的请求数。不需要上 Prometheus 也不需要看 Grafana 图表——一行日志就能快速判断流量分布是否均匀、是否有"冷节点"。

一致性哈希:不同的路由需求

P2C 适合大多数无状态 RPC 调用——哪个节点处理请求都一样。但有些场景需要将同一个 key 的请求发到同一个节点:有状态的分片缓存(同一片数据总是被同一个节点处理)、哈希路由的流量控制(A/B 测试按用户 ID 划分)。

go-zero 提供了一致性哈希作为备选的负载均衡策略:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// zrpc/internal/balancer/consistenthash/consistenthash.go
func (p *picker) Pick(info balancer.PickInfo) (balancer.PickResult, error) {
hashKey := GetHashKey(info.Ctx)
if len(hashKey) == 0 {
return emptyPickResult, status.Error(codes.InvalidArgument,
"[consistent_hash] missing hash key in context")
}

if addrAny, ok := p.hashRing.Get(hashKey); ok {
addr, ok := addrAny.(string)
subConn, ok := p.conns[addr]
return balancer.PickResult{SubConn: subConn}, nil
}

return emptyPickResult, status.Errorf(codes.Unavailable,
"[consistent_hash] no matching conn for hashKey: %s", hashKey)
}

使用时需要在 context 中设置 hash key:

1
2
ctx = consistenthash.SetHashKey(ctx, userId)
client.Call(ctx, req)

一致性哈希使用 go-zero 的 hash.ConsistentHash 实现(默认 100 个虚拟节点),当节点列表变化时只影响相邻节点的数据流,大部分 key 的映射保持不变。与 P2C 不同,一致性哈希不做延迟感知和健康度判断——它只负责将同一个 key 固定映射到同一个地址。如果需要健康管理,需要在上层(如代码层面)自行处理。

Subset:实用的大规模保护

在所有 resolver 的 update 闭包中,你都会看到一句话:

1
vals := subset(sub.Values(), subsetSize)

subsetSize 固定为 32。如果某个服务有 500 个实例,客户端只需要从其中随机挑选 32 个来建立连接,而不是 500 个。这不是因为 P2C 算法不能处理 500 个节点——而是因为 gRPC 为每个地址创建 SubConn 的实际资源开销不容忽视(每个 SubConn 背后是一个 HTTP/2 连接和多个 stream)。

1
2
3
4
5
6
7
8
9
func subset(set []string, sub int) []string {
rand.Shuffle(len(set), func(i, j int) {
set[i], set[j] = set[j], set[i]
})
if len(set) <= sub {
return set
}
return set[:sub]
}

随机打乱 + 截取前 32 个——每个客户端持有的实例子集不同但覆盖不重叠。最终结果是每个实例的客户端连接数大致均等(因为每个实例进入不同客户端子集的概率相同),不存在"特定几个实例被所有客户端连接导致负载不均"的问题。在生产环境中,32 这个值在大多数场景下足够——500 个实例本身已经过于冗余,32 个实例的子集能涵盖足够多的故障转移候选。

完整链路串联

把前文分散在各处的代码片段拼成一张完整的数据流图:

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
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
┌─────────────────────────────────────────────────────────┐
│ 服务端 │
│ RpcPubServer.Start() │
│ → registerEtcd() │
│ → discov.NewPublisher(hosts, key, listenOn) │
│ → Registry.GetConn(hosts) // 复用 etcd 连接 │
│ → cli.Grant(TTL=10s) // 创建租约 │
│ → cli.Put(key, value, WithLease) │
│ → keepAliveAsync(cli) │
│ ├─ KeepAlive channel // 心跳续约 │
│ ├─ Watch channel // 监听自己 key 被删 │
│ └─ Pause/Resume/Quit // 生命周期控制 │
└──────────────┬──────────────────────────────────────────┘
│ etcd 存储
│ /greet.rpc/leaseID → "192.168.1.10:8080"

┌─────────────────────────────────────────────────────────┐
│ 客户端 │
│ NewClient → BuildTarget │
│ "discov:///etcd:2379?key=greet.rpc" │
│ → grpc.DialContext → resolver.Build │
│ → discov.NewSubscriber(hosts, key) │
│ → Registry.Monitor(key, container) │
│ → cluster.getClient() // 复用 etcd 连接 │
│ → cluster.load() // 全量 Get + 初始填充 │
│ → cluster.watch() // Watch 变更 │
│ ├─ watchStream: 事件循环 │
│ │ └─ handleWatchEvents │
│ │ → UpdateListener.OnAdd/OnDelete │
│ │ → container.addKv/removeKey │
│ │ → notifyChange() │
│ └─ watchConnState: 连接断开恢复 │
│ → cluster.reload() │
│ → sub.AddListener(update) // 注册 update 闭包 │
│ → update(){ cc.UpdateState(addrs) } │
└──────────────┬──────────────────────────────────────────┘
│ gRPC 接收新地址列表

┌─────────────────────────────────────────────────────────┐
│ gRPC Balancer │
│ base.NewBalancerBuilder("p2c_ewma") │
│ → 为每个 Address 建立 SubConn │
│ → 监听 SubConn 连接状态变化 │
│ → 仅 Ready SubConn 进入 picker │
│ → 每次 RPC 调用: Picker.Pick │
│ ├─ 从健康 Ready 节点中随机选 2 个 │
│ ├─ choose(c1, c2): 比较 load() │
│ │ load = sqrt(lag) × (inflight + 1) │
│ │ ├─ 正常: 选 load 较小的 │
│ │ └─ 强制拾取: 超过 1s 未选中的节点直接选中 │
│ ├─ inflight++ │
│ └─ 返回 PickResult{SubConn, Done} │
│ │
│ RPC 调用完成 → Done(info) │
│ ├─ inflight-- │
│ ├─ EWMA 更新 lag: │
│ │ w = exp(-td / 10s) │
│ │ new_lag = old_lag × w + current_lag × (1-w) │
│ └─ EWMA 更新 success: │
│ new_success = old_success × w + (0 or 1000) × (1-w) │
│ (根据 codes.Acceptable 判断) │
└─────────────────────────────────────────────────────────┘

从服务端"我来也"的一声宣告,到客户端精准地在一个最优的节点上执行请求,整条链路涉及了 4 个模块、十余个类型、3 个并发的 goroutine 循环,但对业务代码来说,这一切只浓缩为两行配置和一行调用:

1
2
3
4
5
6
7
8
// 服务端配置
RpcServerConf:
Etcd:
Hosts: ["127.0.0.1:2379"]
Key: "greet.rpc"

// 客户端调用
client.SayHello(ctx, &req)

这就是框架的真正价值——基础设施的复杂性与业务代码的简洁性之间,隔着一套精密的自动化管道。 go-zero 的服务发现和负载均衡就是这条管道的关键一段。

总结

本文从服务端注册到客户端选端,完整走完了 go-zero 的服务发现与负载均衡主链路:

服务注册遵循"声明式配置 + 自动 KeepAlive"的模式。 你只需在配置中写 etcd 地址和 key 前缀,框架自动处理租约创建、心跳续约、异常重建和优雅退出。Pause/Resume 机制为运行时的流量摘除提供了精确的控制点。

服务发现的核心是"Registry + cluster + UpdateListener"三层抽象。 Registry 管理 etcd 连接的复用(同一端点共享一个连接),cluster 管理 key 级别的 watch 生命周期(不同 key 共享连接但独立 watch),UpdateListener 将 etcd 的增量变化转化为通用的 OnAdd/OnDelete 通知。StateWatcher 和 revision 压缩处理保证了连接断开和 etcd 运维期间的自愈能力。

Resolver 是"服务发现"和"负载均衡"之间的桥梁。 它实现了 gRPC 的 resolver.Builder 接口,将 Subscriber 的变化通过 cc.UpdateState 自动同步到 gRPC 的地址管理层面。四种 scheme(direct、discov、etcd、k8s)覆盖了从开发环境到生产 Kubernetes 的各种部署场景。Subset 机制在所有 resolver 中统一启用,防止大实例集下的连接数爆炸。

P2C + EWMA 的组合让负载均衡从"轮询"升级为"感知"。 两次随机选择将决策复杂度降到 O(1),EWMA 延迟和成功率让算法能够区分快节点和慢节点,强制拾取打破了惯性防止信号腐化。相比简单的轮询或最少连接数算法,P2C 在节点性能异构的场景下显著减少了请求的长尾延迟。

一致性哈希为有状态路由提供了另一种维度。 P2C 适合"任意节点都能处理请求"的无状态场景,而一致性哈希保证同一个 key 总是落到同一个节点——适合分片缓存、A/B 测试等需要路由亲和性的需求。

理解了服务发现和负载均衡之后,我们在下一篇将进入 go-zero 的另一块核心竞争力——弹性保护的三层武器:熔断、降载与限流。当你知道了客户端如何发现和选择节点之后,再看到"节点连续失败导致断路器跳闸,CPU 过载触发降载拒绝请求"的全过程,才能真正理解 go-zero 的弹性体系是怎样从头武装到尾的。