学习目标
学完本章你应该能够:
- 用自己的话讲清熔断、限流、超时、重试四者的分工,面试时能画一张对比表。
- 解释熔断三态模型(Closed / Open / Half-Open)的状态转移条件,并说清"为什么需要半开"。
- 基于 Kitex 的
endpoint.Middleware机制,手写一个带原子计数器的熔断器。 - 对比自定义实现与 Apache Sentinel 两种方案的取舍,知道什么时候该上 Sentinel。
- 指出固定窗口的"临界突发"缺陷,能用滑动窗口改进错误率统计。
前置知识:
- Kitex 基础(建议先看 [[Kitex/中间件]],熔断通过 Middleware 机制嵌入)
- Go 的
sync/atomic基本用法 - 微服务容错三件套:[[Kitex/限流]]、[[Kitex/超时]]、[[Kitex/重试]]
本章你会动手做的事:
- 跑通一个自定义熔断 Middleware,并故意把下游打挂,观察状态从 Closed → Open → Half-Open → Closed。
- 把熔断参数(错误率、最小请求数)调到极端值,体会"误熔断"和"不起作用"。
- 给熔断加上 Fallback 降级,对比"快速失败"和"返回兜底数据"的体验差异。
概述
熔断(Circuit Breaker) 是一种容错机制,用于防止服务雪崩。当某个下游服务持续失败时,熔断器自动切断对该服务的调用,快速失败,给下游服务恢复的时间;之后自动尝试恢复,形成"闭→开→半开→闭/开"的状态循环。
类比:熔断思想源自电路保护——电流过大时保险丝断开,切断电路保护设备。微服务里的"保险丝"就是熔断器:下游服务像一台过载的电器,如果不及时"拉闸",上游会一直把请求灌进去,最终自己也跟着被拖死(雪崩)。拉闸后请求立刻失败(而不是傻等超时),下游才能喘口气恢复。
Kitex 本身没有内置熔断器实现(不像限流有 pkg/limit),但提供了灵活的中间件机制和 rpcinfo 上下文,用户可以轻松构建熔断逻辑。生产环境中通常通过以下三种方式实现:
| 方案 | 适用场景 | 复杂度 |
|---|---|---|
| 自定义 Middleware | 轻量级、简单项目 | ⭐⭐ |
| Apache Sentinel | 需要统一管控台、动态规则 | ⭐⭐⭐ |
| Resilience4j 风格第三方库 | 需要成熟熔断器实现 | ⭐⭐⭐ |
熔断三态模型
熔断器的核心是一个状态机。理解这三个状态及其转移条件,是看懂后面所有代码的前提。
stateDiagram-v2
[*] --> Closed
Closed --> Open : 错误率超过阈值
(且请求数 >= MinRequests)
Open --> HalfOpen : 等待 PassDuration 超时
尝试恢复
HalfOpen --> Closed : 试探请求成功数
达到 SuccessThreshold
HalfOpen --> Open : 试探请求失败
重新打开| 状态 | 含义 | 行为 |
|---|---|---|
| Closed(闭合) | 正常状态,请求通过 | 统计失败率,超过阈值则触发熔断 |
| Open(打开) | 熔断状态,拒绝请求 | 所有请求直接快速失败,不转发到下游 |
| Half-Open(半开) | 试探恢复 | 放行少量试探请求,成功则恢复 Closed,失败则重新 Open |
⚠️ 新手必踩的坑:把"半开"当万能药。如果下游真的挂了,半开试探只会增加无效请求、进一步消耗资源。所以
PassDuration(恢复等待时间)一定要给下游留足恢复窗口,别设太短。
方案一:自定义 Middleware 实现熔断
这是最灵活的方式,利用 Kitex 的 endpoint.Middleware 机制在客户端实现。整体调用链是:熔断检查 → 执行真实调用 → 根据结果记录成功/失败,步骤如下面时序图所示。
sequenceDiagram
participant C as 调用方
participant M as 熔断中间件
participant CB as CircuitBreaker
participant N as 下游 Endpoint
C->>M: 请求
M->>CB: AllowRequest()
alt 熔断已打开
CB-->>M: false
M-->>C: 快速失败(CircuitBreakerOpenError)
else 允许放行
CB-->>M: true
M->>N: next(ctx, req, resp)
N-->>M: err
alt err != nil
M->>CB: RecordFailure()
else 成功
M->>CB: RecordSuccess()
end
M-->>C: 返回结果
end核心结构定义
package circuitbreaker
import (
"context"
"sync/atomic"
"time"
"github.com/cloudwego/kitex/pkg/endpoint"
"github.com/cloudwego/kitex/pkg/klog"
)
// State 熔断器状态
type State int
const (
StateClosed State = 0 // 闭合:正常
StateOpen State = 1 // 打开:熔断
StateHalfOpen State = 2 // 半开:试探
)
// Config 熔断器配置
type Config struct {
ErrorRatio float64 // 错误率阈值,超过此值触发熔断(如 0.5 = 50%)
MinRequests int // 最小请求数,低于此数不触发熔断
PassDuration time.Duration // 半开后等待多久放试探请求(如 30s)
TripeRequests int // 半开状态下允许的试探请求数(如 3)
SuccessThreshold int // 半开状态下需要多少成功才算恢复
}
// CircuitBreaker 熔断器
type CircuitBreaker struct {
config Config
state atomic.Int64 // State
failedCount atomic.Int64 // 当前窗口失败次数
successCount atomic.Int64 // 当前窗口成功次数
totalCount atomic.Int64 // 当前窗口总请求数
lastFailureTime atomic.Int64 // 最后一次失败的时间戳 (Unix nano)
halfOpenSuccess atomic.Int64 // 半开状态下的连续成功次数
halfOpenWindow time.Time // 半开窗口开启时间
mu sync.Mutex // 状态切换保护锁
}
状态判断与切换逻辑
// 步骤 1:根据当前状态决定是否放行请求
func (cb *CircuitBreaker) AllowRequest() bool {
currentState := State(cb.state.Load())
switch currentState {
case StateClosed:
// 闭合状态:放行请求,统计结果
cb.totalCount.Add(1)
return true
case StateOpen:
// 步骤 2:打开状态,检查是否过了恢复等待期
lastFail := time.UnixNano(cb.lastFailureTime.Load())
if time.Since(lastFail) > cb.config.PassDuration {
// 切换到半开
cb.state.Store(int64(StateHalfOpen))
cb.halfOpenSuccess.Store(0)
klog.Warnf("circuit breaker: OPEN → HALF_OPEN, allowing trial requests")
return true
}
// 仍在熔断期,快速拒绝
return false
case StateHalfOpen:
// 步骤 3:半开状态,只允许有限的试探请求
return cb.totalCount.Load()-cb.halfOpenSuccess.Load() <= int64(cb.config.TripeRequests)
default:
return false
}
}
// 步骤 4:记录成功,半开态下达到阈值则恢复闭合
func (cb *CircuitBreaker) RecordSuccess() {
currentState := State(cb.state.Load())
cb.totalCount.Add(1)
cb.successCount.Add(1)
if currentState == StateHalfOpen {
cb.halfOpenSuccess.Add(1)
// 半开状态下达到成功阈值 → 恢复闭合
if cb.halfOpenSuccess.Load() >= int64(cb.config.SuccessThreshold) {
cb.state.Store(int64(StateClosed))
cb.resetCounters()
klog.Info("circuit breaker: HALF_OPEN → CLOSED, service recovered")
}
}
}
// 步骤 5:记录失败,半开态失败则重开,闭合态超阈值则熔断
func (cb *CircuitBreaker) RecordFailure() {
currentState := State(cb.state.Load())
cb.totalCount.Add(1)
cb.failedCount.Add(1)
cb.lastFailureTime.Store(time.Now().UnixNano())
if currentState == StateHalfOpen {
// 半开状态下失败 → 重新打开
cb.state.Store(int64(StateOpen))
klog.Warn("circuit breaker: HALF_OPEN → OPEN, trial failed, re-opening")
return
}
// 闭合状态下检查是否触发熔断
if cb.shouldTrip() {
cb.state.Store(int64(StateOpen))
klog.Warnf("circuit breaker: CLOSED → OPEN, error ratio %.2f exceeded threshold %.2f",
float64(cb.failedCount.Load())/float64(cb.totalCount.Load()),
cb.config.ErrorRatio)
}
}
// shouldTrip 判断是否触发熔断:样本不足 + 错误率超阈值
func (cb *CircuitBreaker) shouldTrip() bool {
total := cb.totalCount.Load()
if total < int64(cb.config.MinRequests) {
return false
}
errorRatio := float64(cb.failedCount.Load()) / float64(total)
return errorRatio >= cb.config.ErrorRatio
}
func (cb *CircuitBreaker) resetCounters() {
cb.failedCount.Store(0)
cb.successCount.Store(0)
cb.totalCount.Store(0)
}
组装为 Middleware
// 步骤 1:检查是否允许请求(熔断打开则快速失败)
func CircuitBreakerMW(cb *CircuitBreaker) endpoint.Middleware {
return func(next endpoint.Endpoint) endpoint.Endpoint {
return func(ctx context.Context, req, resp interface{}) (err error) {
// 步骤 2:熔断检查
if !cb.AllowRequest() {
klog.Warn("circuit breaker: OPEN, request rejected (fast-fail)")
return &CircuitBreakerOpenError{msg: "circuit breaker open"}
}
// 步骤 3:执行实际调用
err = next(ctx, req, resp)
// 步骤 4:根据结果记录成功或失败
if err != nil {
cb.RecordFailure()
} else {
cb.RecordSuccess()
}
return err
}
}
}
// CircuitBreakerOpenError 熔断打开时的错误
type CircuitBreakerOpenError struct{ msg string }
func (e *CircuitBreakerOpenError) Error() string { return e.msg }
func (e *CircuitBreakerOpenError) Is(target error) bool {
_, ok := target.(*CircuitBreakerOpenError)
return ok
}
使用方式
// 步骤 1:创建熔断器
cb := &CircuitBreaker{
config: Config{
ErrorRatio: 0.5, // 错误率超过 50% 触发
MinRequests: 20, // 至少 20 个请求才判定
PassDuration: 30 * time.Second, // 30s 后进入半开
TripeRequests: 3, // 半开允许 3 个试探
SuccessThreshold: 2, // 2 个成功即恢复
},
}
// 步骤 2:客户端注入中间件
client, err := math.NewClient(
"dqq.math",
client.WithMiddleware(CircuitBreakerMW(cb)),
client.WithRPCTimeout(200*time.Millisecond),
)
方案二:集成 Apache Sentinel
Sentinel 是阿里巴巴开源的流量治理组件,提供熔断降级、流量整形、系统自适应保护等能力。当你的熔断规则需要在管控台动态调、需要统一监控时,自己手写就不划算了,直接上 Sentinel。
添加依赖
go get github.com/alibaba/sentinel-golang/pkg/core/config
go get github.com/alibaba/sentinel-golang/pkg adapter
配置熔断规则
import (
"github.com/alibaba/sentinel-golang/core/circuitbreaker"
"github.com/alibaba/sentinel-golang/core/system"
"github.com/alibaba/sentinel-golang/adapters/kitex"
)
// 步骤 1:初始化 Sentinel,指定熔断策略
cfg := config.NewDefaultConfig()
cfg.Sentinel.CircuitBreakerStrategy = circuitbreaker.StrategyErrorCountRatio // 或 ErrorRatio
config.ApplyConfig(cfg)
// 步骤 2:定义熔断规则
rule := circuitbreaker.Rule{
Resource: "dqq.math/Service.Sub",
Strategy: circuitbreaker.StrategyErrorRatio,
RequestCount: 20, // 统计窗口内的请求数
ErrorRatioThreshold: 0.5, // 错误率 50%
MinRequestNumber: 10, // 最小请求数
SlowRatioThreshold: 0.0, // 慢调用比例(可选)
MaxAllowedSlows: 0, // 半开期间最大慢调用数
StatIntervalMs: 10000, // 统计窗口 10s
Temperature: 0, // 温度因子
RecoveryTimeoutMs: 30000, // 恢复超时 30s
VarianceCoefficient: 1.0, // 方差系数
MeanSlowRequestDuration: 0, // 慢调用阈值
}
// 步骤 3:注册规则
circuitbreaker.PutRule(rule)
Sentinel 支持的三种熔断策略
| 策略 | 说明 | 适用场景 |
|---|---|---|
| ErrorCount | 单位时间内错误数超过阈值 | 错误绝对数量可控的场景 |
| ErrorRatio | 错误率超过阈值 | 按比例熔断,更灵活 |
| SlowRatio | 慢调用比例超过阈值 | 应对响应变慢但不报错的情况 |
在 Kitex 中使用
import "github.com/alibaba/sentinel-golang/adapters/kitex"
// 只需注入 Sentinel 提供的中间件即可,无需自己维护状态机
client, err := math.NewClient(
"dqq.math",
client.WithMiddleware(kitex.ClientSentinelMiddleware()),
client.WithRPCTimeout(200*time.Millisecond),
)
Sentinel Dashboard 可以实时查看熔断状态和动态修改规则:
┌──────────────────────────────────────────────┐
│ Sentinel Dashboard │
│ │
│ 资源: dqq.math/Service.Sub │
│ ┌─────────┬─────────┬─────────┐ │
│ │ QPS │ 错误率 │ 慢调用率 │ │
│ │ 1200 │ 12.3% │ 0.5% │ │
│ └─────────┴─────────┴─────────┘ │
│ │
│ 熔断规则: │
│ 策略: ErrorRatio 阈值: 0.5 │
│ 当前状态: ● CLOSED │
│ [编辑规则] [查看监控] │
└──────────────────────────────────────────────┘
方案三:基于滑动窗口的改进实现
固定窗口存在"临界突发"问题(两个窗口边界各 burst 一倍流量)。滑动窗口更平滑:
固定窗口:
|██████████|██████████|██████████|
t0 t1 t2 t3
↑ 边界处可 burst 2x 流量
滑动窗口:
────────────────────────────────→ t
[████████████████████████████]
←────── 最近 N 秒,平滑统计 ──────→
下面用多个桶模拟滑动窗口,把最近 N 秒切成 bucketCount 个时间片,统计时只累加这些桶。
// 步骤 1:定义滑动窗口结构(多个原子计数器桶)
type SlidingWindowCounter struct {
interval time.Duration
buckets []atomic.Int64
bucketCount int
currentIndex int
}
// 步骤 2:初始化所有桶为 0
func NewSlidingWindowCounter(interval time.Duration, buckets int) *SlidingWindowCounter {
arr := make([]atomic.Int64, buckets)
for i := range arr {
arr[i].Store(0)
}
return &SlidingWindowCounter{
interval: interval,
buckets: arr,
bucketCount: buckets,
}
}
// 步骤 3:根据当前时间戳定位桶下标(取模实现环形复用)
func (sw *SlidingWindowCounter) currentBucketIndex() int {
return int(time.Now().UnixNano() / sw.interval.Nanoseconds()) % sw.bucketCount
}
// 步骤 4:计数。成功 +1,失败 -1(负数表达失败,便于统一统计)
func (sw *SlidingWindowCounter) Increment(isSuccess bool) {
idx := sw.currentBucketIndex()
if isSuccess {
sw.buckets[idx].Add(1)
} else {
sw.buckets[idx].Add(-1) // 负数表示失败
}
}
// 步骤 5:遍历所有桶计算错误率
func (sw *SlidingWindowCounter) ErrorRatio() (float64, int) {
var total, failed int64
for _, b := range sw.buckets {
v := b.Load()
if v > 0 {
total += v
} else {
failed += -v
total += v // v 是负数
}
}
if total == 0 {
return 0, 0
}
return float64(failed) / float64(total), int(total)
}
⚠️ 新手必踩的坑:滑动窗口的"桶环形复用"有并发与精度问题。上面是教学简化版;生产里要么用带时间戳的桶(过期即清零),要么用 Sentinel/Hystrix 这类成熟实现。自己造轮子统计时,别在
currentBucketIndex跨桶边界时漏算正在被覆写的旧桶。
熔断 vs 限流 vs 超时 vs 重试
这四者经常配合使用,分工明确。一句话区分:超时管"等多久"、限流管"别压垮自己"、熔断管"别被下游拖死"、重试管"临时失败再试一次"。
| 机制 | 层级 | 目的 | 触发条件 | 行为 |
|---|---|---|---|---|
| 超时 | 客户端 | 防止等待过久 | 超过指定时间未返回 | 放弃等待,返回错误 |
| 限流 | 服务端 | 保护自身不被打垮 | 超过 QPS/并发上限 | 拒绝新请求 |
| 熔断 | 客户端 | 防止下游故障扩散 | 错误率/慢调用率过高 | 切断对下游的调用 |
| 重试 | 客户端 | 提高成功率 | 临时性失败(网络抖动等) | 再次发起请求 |
一次完整请求里四者的配合顺序,看下面这张流程图(熔断是"第一道闸",最该优先判断):
flowchart TD
A[客户端发起请求] --> B{熔断检查}
B -->|熔断打开| BF[快速失败 不浪费资源]
B -->|允许放行| C[超时控制 设定 deadline]
C --> D[负载均衡 选一个实例]
D --> E[发送请求 等待响应]
E --> F{成功?}
F -->|成功| G[记录成功 返回结果]
F -->|失败| H{还有重试次数?}
H -->|是| E
H -->|否| I[返回错误]
I --> J[熔断器记录失败]最佳实践
1. 熔断参数调优
Config{
// 错误率阈值:建议 0.4 ~ 0.6
// 太低容易误熔断,太高起不到保护作用
ErrorRatio: 0.5,
// 最小请求数:建议 10 ~ 30
// 避免样本太少时误判
MinRequests: 20,
// 恢复等待时间:建议 15s ~ 60s
// 给下游足够的恢复时间
PassDuration: 30 * time.Second,
// 半开试探请求数:建议 2 ~ 5
// 太少无法判断,太多可能压垮正在恢复的下游
TripeRequests: 3,
// 半开成功阈值:建议 50% ~ 80% 的试探请求
SuccessThreshold: 2,
}
2. 熔断 + 降级(Fallback)
熔断不是终点,配合降级逻辑提供更友好的用户体验:
// 步骤 1:熔断打开时不再快速失败,而是走降级逻辑
func CircuitBreakerWithFallback(cb *CircuitBreaker) endpoint.Middleware {
return func(next endpoint.Endpoint) endpoint.Endpoint {
return func(ctx context.Context, req, resp interface{}) (err error) {
if !cb.AllowRequest() {
// 熔断打开时执行降级逻辑
return doFallback(req)
}
// 步骤 2:正常路径,仍要记录结果
err = next(ctx, req, resp)
if err != nil {
cb.RecordFailure()
} else {
cb.RecordSuccess()
}
return err
}
}
}
// 步骤 3:降级返回缓存数据、默认值或友好提示
func doFallback(req interface{}) error {
// 返回缓存数据、默认值、或友好提示
return &FallbackError{msg: "服务暂时不可用,请稍后重试"}
}
3. 按方法粒度熔断
不同方法的容错能力不同,应该独立熔断:
// 数学计算服务:Sub 可以熔断,Health 不应该熔断
cbMap := map[string]*CircuitBreaker{
"Sub": {config: Config{ErrorRatio: 0.5, MinRequests: 20}},
"ComplexCalc": {config: Config{ErrorRatio: 0.3, MinRequests: 10}},
}
func PerMethodCBMW(cbMap map[string]*CircuitBreaker) endpoint.Middleware {
return func(next endpoint.Endpoint) endpoint.Endpoint {
return func(ctx context.Context, req, resp interface{}) (err error) {
method := getMethodName(ctx) // 从 ctx 提取方法名
cb, ok := cbMap[method]
if !ok {
return next(ctx, req, resp)
}
// ... 同上的熔断逻辑
}
}
}
4. 监控与告警
熔断状态变化本身就是重要的运维信号:
// 在状态切换时上报指标
func (cb *CircuitBreaker) transition(from, to State) {
klog.Warnf("[CIRCUIT_BREAKER] service=%s state=%s→%s",
serviceName, from, to)
// 上报到 Prometheus / StatsD
metrics.Counter("circuit_breaker_state_change",
prometheus.Labels{"service": serviceName, "from": from.String(), "to": to.String()},
1)
}
注意事项
- 不要对幂等性差的操作过度重试:熔断 + 重试组合时,注意重试可能放大流量
- 熔断是客户端行为:每个客户端独立维护熔断状态,可能出现部分客户端已恢复而另一部分仍在熔断的情况
- 半开不是万能的:如果下游真的挂了,半开试探只会增加无效请求,设置合理的
PassDuration很关键 - 区分错误类型:不是所有错误都应该计入熔断(如参数校验错误不应触发熔断),建议在
RecordFailure前过滤 - 连接池影响:熔断器关闭后,如果连接池也断了,第一个请求仍然可能超时,建议配合 [[预热]] 使用
⚠️ 新手必踩的坑:把所有错误都计入熔断。
RecordFailure前一定要区分错误类型——参数校验错误、context.Canceled这类不该算作"下游故障"。否则一次上游参数写错,就会把健康的下游误熔断。
自测题与动手练习
自测题(合上书能答出来,才算懂):
- 熔断三态 Closed / Open / Half-Open 各自代表什么?从 Closed 切到 Open 的条件是什么?从 Open 切到 Half-Open 的条件又是什么?
- 熔断、限流、超时、重试分别工作在客户端还是服务端?各自"防的是什么"?
- 为什么需要
MinRequests?如果把MinRequests设成 1 会有什么问题? - 固定窗口统计错误率有什么缺陷?滑动窗口是怎么缓解"临界突发"的?
- 自定义熔断 Middleware 里,为什么
AllowRequest和RecordSuccess/RecordFailure都要用atomic操作?
动手练习(建议真做一遍):
- 把方案一的状态机跑起来,用
go test或压测工具把下游打挂,观察日志里CLOSED → OPEN → HALF_OPEN → CLOSED的完整转移。 - 故意把
ErrorRatio调到 0.99、MinRequests调到 1,体会"熔断根本不触发";再反向调到 0.01,体会"稍微错一个就熔断"的误伤。 - 给方案一的熔断加上
CircuitBreakerWithFallback,让熔断打开时返回一个本地缓存的默认值,对比"快速失败报错"和"降级兜底"两种用户体验。
本章小结
- 熔断是客户端的"保险丝":下游持续失败时自动切断调用、快速失败,给下游恢复窗口,之后再自动试探恢复。
- 三态模型是核心:Closed(正常统计)→ Open(熔断拒绝)→ Half-Open(试探恢复),转移条件由错误率、
MinRequests、PassDuration、SuccessThreshold共同决定。 - Kitex 无内置熔断,但 Middleware 机制让自定义实现很轻量;需要管控台与动态规则时优先选 Apache Sentinel。
- 固定窗口有临界突发缺陷,滑动窗口统计更平滑,但生产环境建议直接用成熟库而非手搓。
- 四件套(超时/限流/熔断/重试)要组合使用:熔断是第一道闸,先判断能不能调,再谈超时、负载均衡、重试。
下一篇可以顺着看 [[Kitex/限流]]、[[Kitex/超时]]、[[Kitex/重试]],把容错四件套补齐。
相关笔记
- [[Kitex/限流]] — 服务端限流,与熔断互补
- [[Kitex/超时]] — 熔断的前提配置,没有超时熔断无从谈起
- [[Kitex/重试]] — 熔断 + 重试的组合策略
- [[Kitex/负载均衡]] — 负载均衡 + 熔断可实现自动摘除故障节点
- [[Kitex/中间件]] — 熔断通过 Middleware 机制嵌入