为什么你的服务需要熔断器
你的订单服务调用支付服务,某天支付服务突然变慢——从50ms飙到5秒。你的订单服务每个请求都在等5秒,线程池迅速耗尽,整个服务雪崩。最后,一个下游服务的慢响应,把你整个系统拖垮了。
这就是级联故障。熔断器(Circuit Breaker)是解决这个问题的经典模式:当下游错误率超过阈值时,主动切断调用,快速失败,给下游喘息恢复的时间。
熔断器的三态机制
熔断器有三个状态,像一个智能开关:
- Closed(关闭):正常放行所有请求,同时统计失败率
- Open(打开):熔断状态,所有请求直接失败,不发出实际调用
- Half-Open(半开):试探性放行少量请求,根据结果决定恢复或重新熔断
状态流转:Closed →(失败率超阈值)→ Open →(冷却时间到期)→ Half-Open →(探测成功)→ Closed,或 →(探测失败)→ Open
用Go实现一个生产级熔断器
核心数据结构
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
|
package breaker
import (
"errors"
"sync"
"time"
)
type State int
const (
StateClosed State = iota // 关闭:正常放行
StateOpen // 打开:快速失败
StateHalfOpen // 半开:试探放行
)
var ErrCircuitOpen = errors.New("circuit breaker is open")
type Config struct {
RequestThreshold int // 触发熔断的最小请求量
FailureRatio float64 // 失败率阈值(0~1)
OpenTimeout time.Duration // 熔断持续时间
HalfOpenMaxCalls int // 半开状态最大探测请求数
HalfOpenSuccessNeed int // 半开状态恢复所需连续成功数
}
type CircuitBreaker struct {
mu sync.Mutex
state State
config Config
window *slidingWindow
openedAt time.Time
halfOpenSuccess int
halfOpenCalls int
}
func New(cfg Config) *CircuitBreaker {
return &CircuitBreaker{
state: StateClosed,
config: cfg,
window: newSlidingWindow(10 * time.Second),
}
}
|
滑动窗口统计
用滑动窗口记录最近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
39
40
41
42
43
44
45
46
47
48
49
|
type slidingWindow struct {
mu sync.Mutex
duration time.Duration
buckets []bucket
}
type bucket struct {
timestamp time.Time
success int
failure int
}
func newSlidingWindow(d time.Duration) *slidingWindow {
return &slidingWindow{duration: d}
}
func (w *slidingWindow) Record(success bool) {
w.mu.Lock()
defer w.mu.Unlock()
now := time.Now()
// 清理过期桶
cutoff := now.Add(-w.duration)
alive := w.buckets[:0]
for _, b := range w.buckets {
if b.timestamp.After(cutoff) {
alive = append(alive, b)
}
}
w.buckets = alive
// 追加新记录
if success {
w.buckets = append(w.buckets, bucket{timestamp: now, success: 1})
} else {
w.buckets = append(w.buckets, bucket{timestamp: now, failure: 1})
}
}
func (w *slidingWindow) Stats() (total, failures int) {
w.mu.Lock()
defer w.mu.Unlock()
cutoff := time.Now().Add(-w.duration)
for _, b := range w.buckets {
if b.timestamp.After(cutoff) {
total += b.success + b.failure
failures += b.failure
}
}
return
}
|
状态机:请求放行与结果记录
这是熔断器的核心逻辑,两个方法配合驱动状态转换:
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
|
// Allow 判断请求是否被放行
func (cb *CircuitBreaker) Allow() error {
cb.mu.Lock()
defer cb.mu.Unlock()
switch cb.state {
case StateClosed:
return nil
case StateOpen:
if time.Since(cb.openedAt) >= cb.config.OpenTimeout {
cb.state = StateHalfOpen
cb.halfOpenSuccess = 0
cb.halfOpenCalls = 0
return nil // 放行第一个探测请求
}
return ErrCircuitOpen
case StateHalfOpen:
if cb.halfOpenCalls < cb.config.HalfOpenMaxCalls {
cb.halfOpenCalls++
return nil
}
return ErrCircuitOpen
}
return nil
}
// Record 记录请求结果,驱动状态转换
func (cb *CircuitBreaker) Record(err error) {
cb.mu.Lock()
defer cb.mu.Unlock()
success := err == nil
switch cb.state {
case StateClosed:
cb.window.Record(success)
total, failures := cb.window.Stats()
if total >= cb.config.RequestThreshold {
if float64(failures)/float64(total) >= cb.config.FailureRatio {
cb.open()
}
}
case StateHalfOpen:
if success {
cb.halfOpenSuccess++
if cb.halfOpenSuccess >= cb.config.HalfOpenSuccessNeed {
cb.close()
}
} else {
cb.open()
}
}
}
func (cb *CircuitBreaker) open() {
cb.state = StateOpen
cb.openedAt = time.Now()
}
func (cb *CircuitBreaker) close() {
cb.state = StateClosed
cb.window = newSlidingWindow(10 * time.Second)
}
|
包装调用与使用示例
1
2
3
4
5
6
7
8
9
|
// Do 封装请求执行,自动应用熔断逻辑
func (cb *CircuitBreaker) Do(fn func() error) error {
if err := cb.Allow(); err != nil {
return err
}
err := fn()
cb.Record(err)
return err
}
|
实际使用:
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
|
func main() {
cb := New(Config{
RequestThreshold: 20,
FailureRatio: 0.5,
OpenTimeout: 30 * time.Second,
HalfOpenMaxCalls: 3,
HalfOpenSuccessNeed: 2,
})
callAPI := func() error {
resp, err := http.Get("https://api.example.com/pay")
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 500 {
return errors.New("server error")
}
return nil
}
err := cb.Do(callAPI)
if errors.Is(err, ErrCircuitOpen) {
// 熔断中,走降级逻辑
fmt.Println("支付服务暂时不可用,请稍后重试")
}
}
|
开源库选型
自实现适合理解原理和定制需求。生产环境也可用成熟的开源库:
| 库 |
特点 |
适用场景 |
| sony/gobreaker |
轻量,API简洁,支持状态转换回调 |
标准熔断需求 |
| afex/hystrix-go |
Netflix Hystrix的Go移植 |
功能最全,但已停止维护 |
| go-kratos内熔断器 |
与框架深度集成 |
已用kratos的项目 |
关键配置建议
- RequestThreshold:至少20-50,太小容易误判
- FailureRatio:0.5比较合理,太敏感会导致正常波动也触发熔断
- OpenTimeout:10-60秒,太短下游没来得及恢复,太长影响可用性
- HalfOpenSuccessNeed:2-3次,避免单次偶然成功导致误恢复
熔断器不是银弹
熔断器解决的是快速失败的问题,但不是容错的全部。完整的容错体系还需要:
- 超时:每个调用必须有超时,没有超时就没有熔断的基础
- 重试:瞬态错误用重试解决,但必须配指数退避,否则加剧下游压力
- 降级:熔断后返回缓存数据或默认值,而不是直接给用户报错
- 隔离仓(Bulkhead):不同下游用不同连接池,避免一个下游拖垮全部
正确顺序:超时 → 重试 → 熔断 → 降级。先保证单次调用有超时,再决定是否重试,熔断拦截持续失败,最后降级兜底。