代码路径:go\pkg\mod\github.com\zeromicro\go-zero@v1.10.2\zrpc\internal\balancer\p2c\p2c.go
const (
Name = "p2c_ewma" // balancer注册名字,gRPC识别用
decayTime = int64(time.Second * 10) // EWMA衰减时间 10s
forcePick = int64(time.Second) // 1秒,强制选中保护:一个节点太久没被选中,强制给一次机会
initSuccess = 1000 // 初始化成功计数,节点刚创建默认健康
throttleSuccess = initSuccess / 2 // 低于500,判定不健康
penalty = int64(math.MaxInt32) // 惩罚值,负载极大
pickTimes = 3 // 最多尝试3次随机挑选两个健康节点
logInterval = time.Minute // 每分钟打印一次统计日志
)
p2cPickerBuilder 构建器
func (b *p2cPickerBuilder) Build(info base.PickerBuildInfo) balancer.Picker {
// gRPC resolver 给到的当前可用实例连接列表
readySCs := info.ReadySCs
if len(readySCs) == 0 {
return base.NewErrPicker(balancer.ErrNoSubConnAvailable)
}
conns := make([]*subConn, 0, len(readySCs))
for conn, connInfo := range readySCs {
// 每个实例包装成subConn,保存 lag、inflight、success 等统计指标。
conns = append(conns, &subConn{
addr: connInfo.Address,
conn: conn,
success: initSuccess, // 新节点初始化success=1000,标记健康
})
}
return &p2cPicker{
conns: conns,
r: rand.New(rand.NewSource(time.Now().UnixNano())),
stamp: syncx.NewAtomicDuration(),
}
}
当 etcd 服务发现的 resolver 感知到服务实例发生新增或下线变化时,会触发调用Build方法重建 picker。info.ReadySCs是 gRPC resolver 传递过来的当前可用实例连接集合,遍历该集合,把每一个实例封装成subConn结构体,用于保存实例对应的 EWMA 延迟 lag、正在处理请求数 inflight、健康分数 success 等负载统计指标。
请求进来时P2C 核心选择逻辑
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:
// 实例数量 = 1:直接选这个
chosen = p.choose(p.conns[0], nil)
case 2:
// 实例数量 = 2:直接对比两个
// 对比两者load()负载,选出负载低的
chosen = p.choose(p.conns[0], p.conns[1])
default:
var node1, node2 *subConn
// 最多尝试3次随机取出2个健康节点
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]
// healthy() 判断success >500
if node1.healthy() && node2.healthy() {
break
}
}
chosen = p.choose(node1, node2)
}
// 选中后,inflight+1:正在处理请求数+1
atomic.AddInt64(&chosen.inflight, 1)
atomic.AddInt64(&chosen.requests, 1)
return balancer.PickResult{
SubConn: chosen.conn,
Done: p.buildDoneFunc(chosen), // 回调函数,rpc调用结束触发
}, nil
}
当可用实例数量为 1 时,直接选用该实例;实例数量等于 2 时,直接对比两个实例负载;实例数量大于 2 时,最多尝试 3 次随机挑选,选出两个不同且健康的节点 node1、node2。随后调用 choose 方法对比两者的 load() 综合负载,选择负载更低的实例。选中实例后将其 inflight 并发请求计数加 1,标记该实例新增一个处理中的请求;最后返回选中的连接以及 Done 回调函数,待 RPC 调用完成后,由 gRPC 自动执行该回调更新统计指标。核心特点:每次仅随机挑选 2 个实例对比,不会遍历全部实例;最多重试 3 次,尽量避免选到不健康节点。
choose (c1,c2):对比两个实例负载,选出更好的
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 上一次被选中的时间戳
pick := atomic.LoadInt64(&c2.pick)
// c2 距离上次被选中已经超过 1s
if start-pick > forcePick && atomic.CompareAndSwapInt64(&c2.pick, pick, start) {
return c2
}
atomic.StoreInt64(&c1.pick, start)
return c1
}
P2C 算法在选择节点时,会调用c.load()计算实例综合负载,负载由 EWMA 平滑延迟 lag 与当前并发请求数 inflight 共同算出,正常情况下优先选择负载更小的实例;同时内置饥饿保护机制:当随机选出两个候选节点 c1、c2,c1 负载更低,但若 c2 为健康节点且距离上一次被选中已超过 1 秒,则本次强制选择 c2,仅本次生效并刷新 c2 的选中时间戳,防止健康实例长期分配不到请求而出现饥饿。
rpc 调用结束执行EWMA 更新
func (p *p2cPicker) buildDoneFunc(c *subConn) func(info balancer.DoneInfo) {
// Pick选中节点的时间
start := int64(timex.Now())
return func(info balancer.DoneInfo) {
// 当前正在跑的请求计数 -1
atomic.AddInt64(&c.inflight, -1)
now := timex.Now()
last := atomic.SwapInt64(&c.last, int64(now))
// td:距离上一次更新这个节点统计数据过去了多久(单位 ns)
// td = now - last,last 是上一次更新 lag/success 的时间戳
td := int64(now) - last
if td < 0 {
td = 0
}
// decayTime = 10s(常量,固定 10 秒)
w := math.Exp(float64(-td) / float64(decayTime))
//本次rpc耗时!
lag := int64(now) - start
if lag < 0 {
lag = 0
}
olag := atomic.LoadUint64(&c.lag)
if olag == 0 {
w = 0
}
// EWMA公式:新lag = old_lag * w + new_lag * (1-w)
atomic.StoreUint64(&c.lag, uint64(float64(olag)*w+float64(lag)*(1-w)))
// 更新success健康分数,失败请求success降低
success := initSuccess // 本次请求对应的健康分数
if info.Err != nil && !codes.Acceptable(info.Err) {
success = 0
}
osucc := atomic.LoadUint64(&c.success) //旧的 success 值
// w:衰减权重
atomic.StoreUint64(&c.success, uint64(float64(osucc)*w+float64(success)*(1-w)))
stamp := p.stamp.Load()
if now-stamp >= logInterval {
if p.stamp.CompareAndSwap(stamp, now) {
p.logStats()
}
}
}
}
EWMA 的计算公式为:
newLag = oldLag \times w + currentLag \times (1-w)
其中衰减权重 (w=e^{-td/10s}),td 代表距离上一次更新指标的时间差。节点越久没有请求打到,w 就越小,历史 lag 的影响力随之降低,旧的延迟数据会随时间逐步遗忘;lag 保存在 subConn.lag 中,代表经过平滑后的 EWMA 延迟。框架会用同一套 EWMA 逻辑同步更新 success 健康分数:RPC 调用如果返回链路层面不可接受的错误(不是业务正常错误),本次 success 取值 0;正常请求则取满分。当 success 分数低于阈值,healthy () 返回 false,该节点被标记为不健康,在 Pick 选择实例阶段会被跳过。
subConn 结构体:单个实例的统计数据
type subConn struct {
lag uint64 // EWMA平滑延迟,核心指标
inflight int64 // 当前正在处理的请求数(并发)
success uint64 // 健康分数,判断节点是否健康
requests int64 // 总请求计数
last int64 // 上次更新lag的时间戳
pick int64 // 上次被选中的时间戳(用于forcePick防饥饿)
addr resolver.Address
conn balancer.SubConn
}
综合负载计算(choose 函数使用)
go-zero 的 P2C 负载均衡在计算节点综合负载时,并不是只依靠 EWMA 平滑延迟 lag,而是把延迟和当前正在处理的并发请求数 inflight 结合起来计算,公式为 load = sqrt(lag+1) × (inflight+1)。即便某个实例单次请求延迟很低,如果该实例上堆积了大量正在处理的请求,inflight 数值变大,最终算出的综合负载 load 也会升高。在节点选择时,框架会优先挑选综合负载更小的实例,从而自动减少高并发实例的流量分配,避免单一实例请求堆积,实现延迟与并发双维度的负载感知。
健康实例判断
func (c *subConn) healthy() bool {
return atomic.LoadUint64(&c.success) > throttleSuccess
}
健康判断:success > 500 才认为健康,不健康节点会尽量避开。
init 注册
func init() {
balancer.Register(newBuilder())
}
程序启动自动注册 p2c balancer 到 gRPC,所以业务代码不用配置,自动启用。
EWMA
EWMA = Exponentially Weighted Moving Average,指数加权移动平均。 换句话说:给最近的数据更高权重,久远的数据权重指数衰减,用来平滑统计延迟,反映节点最近的响应快慢。
解决了轮询面对”慢实例”时的热点问题
传统轮询负载均衡会均匀分发请求,无法感知实例响应快慢,一旦集群中出现慢实例,依旧持续分配流量,造成请求堆积、形成热点。例如集群有 A、B、C 三个实例,A、B 正常 10ms 返回,C 是慢实例需要 1000ms 返回;轮询按照 A→B→C 依次分发,每 3 个请求就会有 1 个打到 C,打到 C 的请求需要等待 1 秒,inflight 请求不断堆积,容易大量超时。而 go-zero 的 P2C+EWMA 算法,会统计每个实例的 EWMA 延迟 lag 和当前并发 inflight,计算综合负载。本例中 C 节点 lag 很高,综合负载变大,每次请求随机挑选两个实例对比负载时,C 被选中的概率大幅降低,流量会自动绕开慢实例,解决了轮询无法规避慢实例带来请求堆积的短板。