go-zero熔断器

go-zero 在每个 RPC 客户端和 HTTP 服务中内置了 Google SRE 风格的熔断器——无需手动配置即可自动生效。

雪崩形成

雪崩其实是分布式系统中一种连锁故障的现象,在一个分布式系统中,局部故障是不可避免的。如果不能将局部故障控制好,导致局部故障被正反馈循环,就会出现连锁故障,最终导致整体系统崩溃。一般来说,出现雪崩会经历以下三个阶段

局部服务过载:一切的开端,通常是系统中某个服务(我们称之为服务C)的处理能力出现瓶颈。导致过载的原因多种多样:可能是程序自身的Bug导致性能劣化;可能是流量洪峰超出了服务容量;也可能是部署的机器实例宕机,导致整体处理能力下降。当服务的QPS超过其处理极限时,它会开始表现出响应变慢、内部资源(如内存、线程)消耗加剧等症状。

资源耗尽与服务不可用:随着过载情况的加剧,积压在服务C内部的请求会越来越多。这些积压的请求会持续消耗着服务器的内存、CPU、线程池甚至是文件句柄等宝贵资源。当其中任何一项资源被耗尽时,服务C就会开始大量报错,甚至频繁崩溃,最终对外呈现为“不可用”状态。

当服务C的某个实例崩溃后,上游的负载均衡机制会把原本发往这个实例的流量,自动转发给其它正常的的实例,这样的话,就加速了整个服务C集群的全面瘫痪。

故障沿调用链路逆向蔓延:到这一步,就轮到服务C的上游调用者(服务B)遭殃了。由于服务C无法及时响应,服务B的请求线程会大量阻塞在等待响应上,这同样会快速耗尽服务B自身的线程和资源。很快,服务B也变得过载、不可用。

这个过程会像多米诺骨牌一样,沿着调用链路一路逆向传播(服务B -> 服务A),最终导致整个调用链上的所有服务全部瘫痪,系统发生“雪崩”。

go-zero熔断器

熔断器在滑动窗口内统计错误率。当错误率超过阈值时,熔断器打开,后续请求立即失败(快速失败),不再等待慢速下游。冷却期结束后,允许一个探测请求通过——若成功则重新关闭熔断器。

window            = time.Second * 10   // 统计总窗口 10s
buckets           = 40                 // 40个桶,每个桶250ms
forcePassDuration = time.Second        // 探测放行间隔:1s
k                 = 1.5                // 基准系数
minK              = 1.1                // k最小值,最低不能低于1.1
protection        = 5                  // 保护阈值,请求量很小的时候,不轻易熔断

*collection.RollingWindow [int64, *bucket]

其中int64指每个桶的统计计数(成功、失败、丢弃都是 int64),*bucket是窗口里面每一个 “桶” 的结构体,每个桶保存一段时间内的请求指标:Success、Failure 计数。

func newGoogleBreaker() *googleBreaker {
    // 计算每个桶时长:`10s /40 = 250ms`,创建 40 个桶的滑动窗口。
    bucketDuration := time.Duration(int64(window) / int64(buckets))
    st := collection.NewRollingWindow[int64, *bucket](func() *bucket {
        return new(bucket)
    }, buckets, bucketDuration)
    return &googleBreaker{
        stat:     st,
        k:        k,
        proba:    mathx.NewProba(),
        lastPass: syncx.NewAtomicDuration(),
    }
}

也就是:窗口拆成 40 个小桶,每个桶负责记录 250ms 内的请求统计数据。 随着时间流逝,旧桶会被重置、复用(环形覆盖),只保留最近 10 秒的数据,更早的数据自动丢弃。

type googleBreaker struct {
    k        float64
    stat     *collection.RollingWindow[int64, *bucket] // 滑动窗口,存放40个bucket,记录success/fail/drop
    proba    *mathx.Proba     // 随机概率发生器,用来按dropRatio随机丢弃请求
    lastPass *syncx.AtomicDuration // 原子保存【上一次放行请求的时间戳】,用于兜底探测
}

// history() 聚合统计后的返回结果
type windowResult struct {
    accepts        int64 // 成功请求总数
    total          int64 // 总请求数(成功+失败)
    failingBuckets int64 // 存在失败的桶数量
    workingBuckets int64 // 有成功、无失败的健康桶数量
}

history ():聚合窗口所有桶的数据

func (b *googleBreaker) history() windowResult {
    var result windowResult
    b.stat.Reduce(func(b *bucket) {
        result.accepts += b.Success
        result.total += b.Sum
        if b.Failure > 0 {
            // 只要当前桶有失败,直接清零健康桶计数
            result.workingBuckets = 0
        } else if b.Success > 0 {
            // 没有失败,并且有成功,健康桶+1
            result.workingBuckets++
        }
        if b.Success > 0 {
            // 当前桶有成功,直接清零失败桶计数
            result.failingBuckets = 0
        } else if b.Failure > 0 {
            // 没有成功、有失败,失败桶+1
            result.failingBuckets++
        }
    })
    return result
}

桶 1:Success请求成功数=10,Failure请求失败数=0,桶 2:Success=10,Failure=0,桶 3:Success=10,Failure=0,桶 4:Success=0,Failure=5,遍历顺序是桶 1→桶 2→桶 3→桶 4,遍历桶 1 时,没有失败且存在成功,workingBuckets 赋值为 1,遍历桶 2 时,没有失败且存在成功,workingBuckets 增加至 2,遍历桶 3 时,没有失败且存在成功,workingBuckets 增加至 3,遍历桶 4 时,b.Failure>0,直接将 workingBuckets 置为 0,最终 workingBuckets 等于 0,前面三个健康桶的连续健康计数全部失效。

workingBuckets 并不是用来统计滑动窗口内健康桶的总数量,而是记录最新连续无失败且有成功请求的桶长度,用来衡量下游服务近期是否持续稳定。Reduce 遍历顺序为从旧桶到新桶,只有最新连续一批桶全部请求成功、不存在失败时,workingBuckets 才会持续累加;一旦最新的桶出现失败请求,连续健康的 “连胜记录” 会直接清零,不再对丢弃概率dropRatio做衰减打折。与之对应的failingBuckets同理,记录最新连续失败的桶数量,只要最新桶出现成功,连续失败计数就会清零。

1:算出权重

w = b.k – (b.k-minK)*float64(history.failingBuckets)/buckets 权重计算

如果没有失败桶 w=1.5-0.01*0=1.5;失败桶为10时,w=1.5-0.01*10=1.4;权重下降,失败桶为30,w=1.5-0.01*30=1.2 ,随着失败桶越来越多,w 从 1.5 逐步下降,最低锁死在 1.1,即mathx.AtLeast(w, minK)

2:算出可接受请求数

weightedAccepts := mathx.AtLeast (w, minK) * float64 (history.accepts) 加权可接受请求数,是 SRE 过载公式里的核心变量。

weightedAccepts = 1.5*100 = 150 ,含义是当前成功 100 笔,系统可以承受最多 150 总请求,超过才会开始考虑丢弃。weightedAccepts = 1.3*100 = 150 ,含义是系统已经有一半桶出现失败,健康系数下降,最多只能扛 130 总请求,更容易触发丢弃。失败桶很多,w 算出来低于 1.1,触发 AtLeast 保底,w 不会小于 1.1,加权可接受请求最低是 1.1 * 成功数,防止 w 无限变小导致无限丢弃。

3:算出丢弃概率

dropRatio := (float64(history.total-protection) – weightedAccepts) / float64(history.total+1)

weightedAccepts=150,假设总请求 history.total=160

(dropRatio = (160-5 -150)/(160+1) = 5 /161 ≈0.031)丢弃概率约 3.1%,少量丢包。如果下游恶化weightedAccepts=110,总请求 history.total=160,(dropRatio=(160-5 -110)/(160+1)=45/161≈0.279\)

丢弃概率 27.9%,大量请求被熔断拒绝。

当滑动窗口内总请求数 ≤5时,total-protection ≤0,分子大概率为负,dropRatio≤0,直接放行请求,不会触发熔断。

func (b *googleBreaker) accept() error {
    var w float64
    history := b.history() // 聚合滑动窗口所有桶,拿到成功数、总请求、失败桶、健康桶
    // 动态计算w:存在失败桶越多,w越小,最低保底 minK=1.1
    w = b.k - (b.k-minK)*float64(history.failingBuckets)/buckets
    weightedAccepts := mathx.AtLeast(w, minK) * float64(history.accepts)

    // SRE 原始公式:dropRatio = max(0, (total - k*accepts) / (total +1))
    // 这里做了改动,增加动态w,还有protection保护
    dropRatio := (float64(history.total-protection) - weightedAccepts) / float64(history.total+1)
    if dropRatio <= 0 {
        return nil // 丢弃概率<=0:全部放行,不熔断
    }

    // ========== 兜底探测逻辑 ==========
    lastPass := b.lastPass.Load()
    if lastPass > 0 && timex.Since(lastPass) > forcePassDuration {
        b.lastPass.Set(timex.Now())
        return nil // 距离上次放行超过1s,强制放行这条探测请求
    }

    // 再乘以健康桶权重:健康桶越少,丢弃概率进一步放大
    dropRatio *= float64(buckets-history.workingBuckets) / buckets

    // 按概率随机丢弃:随机命中,则返回熔断错误
    if b.proba.TrueOnProba(dropRatio) {
        return ErrServiceUnavailable
    }
    b.lastPass.Set(timex.Now())
    return nil
}

即使当前丢弃概率很高,只要距离上一次放行请求超过 1 秒,强制放行本次请求作为探测。

b.proba.TrueOnProba(dropRatio):0~1 之间随机数,如果命中概率,返回 ErrServiceUnavailable,拒绝请求。 没有命中概率:放行,更新 lastPass。

统计埋点函数

func (b *googleBreaker) markDrop() {
    b.stat.Add(drop)      // 请求被熔断器直接丢弃
}
func (b *googleBreaker) markFailure() {
    b.stat.Add(fail)      // 请求发出去了,下游返回失败
}
func (b *googleBreaker) markSuccess() {
    b.stat.Add(success)   // 请求成功
}
type bucket struct {
	Success int64
	Failure int64
	Drop    int64
	Sum     int64 // Sum = Success + Failure + Drop
}

调用 stat.Add(event) 的时候,底层会找到当前时间对应的桶,给对应字段 + 1,Sum 同步 + 1。