Contents

熔断器实战:用Go实现优雅的服务降级与自动恢复

为什么你的服务需要熔断器

你的订单服务调用支付服务,某天支付服务突然变慢——从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次,避免单次偶然成功导致误恢复

熔断器不是银弹

熔断器解决的是快速失败的问题,但不是容错的全部。完整的容错体系还需要:

  1. 超时:每个调用必须有超时,没有超时就没有熔断的基础
  2. 重试:瞬态错误用重试解决,但必须配指数退避,否则加剧下游压力
  3. 降级:熔断后返回缓存数据或默认值,而不是直接给用户报错
  4. 隔离仓(Bulkhead):不同下游用不同连接池,避免一个下游拖垮全部

正确顺序:超时 → 重试 → 熔断 → 降级。先保证单次调用有超时,再决定是否重试,熔断拦截持续失败,最后降级兜底。