Kitex熔断

2024-01-17T14:21:02+08:00 | 7分钟阅读 | 更新于 2024-01-17T14:21:02+08:00

@

概述

熔断(Circuit Breaker) 是一种容错机制,用于防止服务雪崩。当某个下游服务持续失败时,熔断器自动切断对该服务的调用,快速失败,给下游服务恢复的时间;之后自动尝试恢复,形成"闭→开→半开→闭/开"的状态循环。

熔断思想源自电路保护——电流过大时保险丝断开,切断电路保护设备。

Kitex 本身没有内置熔断器实现(不像限流有 pkg/limit),但提供了灵活的中间件机制和 rpcinfo 上下文,用户可以轻松构建熔断逻辑。生产环境中通常通过以下三种方式实现:

方案适用场景复杂度
自定义 Middleware轻量级、简单项目⭐⭐
Apache Sentinel需要统一管控台、动态规则⭐⭐⭐
Resilience4j 风格第三方库需要成熟熔断器实现⭐⭐⭐

熔断三态模型

        错误率超过阈值
    ┌───────┐  ┌──────┐
    │  闭   │──│  开  │
    │ Closed│  │ Open │
    └───────┘  └──────┘
        ▲           │
        │           │ 等待恢复超时
        │      ┌────┴────┐
        └──────│ 半开 Half-Open │
               └─────────────┘
              试探成功 → 闭合
              试探失败 → 重新打开
状态含义行为
Closed(闭合)正常状态,请求通过统计失败率,超过阈值则触发熔断
Open(打开)熔断状态,拒绝请求所有请求直接快速失败,不转发到下游
Half-Open(半开)试探恢复放行少量试探请求,成功则恢复 Closed,失败则重新 Open

方案一:自定义 Middleware 实现熔断

这是最灵活的方式,利用 Kitex 的 endpoint.Middleware 机制在客户端实现。

核心结构定义

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   // 状态切换保护锁
}

状态判断与切换逻辑

func (cb *CircuitBreaker) AllowRequest() bool {
	currentState := State(cb.state.Load())

	switch currentState {
	case StateClosed:
		// 闭合状态:放行请求,统计结果
		cb.totalCount.Add(1)
		return true

	case StateOpen:
		// 打开状态:检查是否过了恢复等待期
		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:
		// 半开状态:只允许有限的试探请求
		return cb.totalCount.Load()-cb.halfOpenSuccess.Load() <= int64(cb.config.TripeRequests)

	default:
		return false
	}
}

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")
		}
	}
}

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)
	}
}

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

func CircuitBreakerMW(cb *CircuitBreaker) endpoint.Middleware {
	return func(next endpoint.Endpoint) endpoint.Endpoint {
		return func(ctx context.Context, req, resp interface{}) (err error) {
			// 1. 检查是否允许请求
			if !cb.AllowRequest() {
				klog.Warn("circuit breaker: OPEN, request rejected (fast-fail)")
				return &CircuitBreakerOpenError{msg: "circuit breaker open"}
			}

			// 2. 执行实际调用
			err = next(ctx, req, resp)

			// 3. 根据结果记录
			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
}

使用方式

// 创建熔断器
cb := &CircuitBreaker{
	config: Config{
		ErrorRatio:       0.5,       // 错误率超过 50% 触发
		MinRequests:      20,        // 至少 20 个请求才判定
		PassDuration:     30 * time.Second, // 30s 后进入半开
		TripeRequests:    3,         // 半开允许 3 个试探
		SuccessThreshold: 2,         // 2 个成功即恢复
	},
}

// 客户端注入
client, err := math.NewClient(
	"dqq.math",
	client.WithMiddleware(CircuitBreakerMW(cb)),
	client.WithRPCTimeout(200*time.Millisecond),
)

方案二:集成 Apache 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"
)

// 初始化 Sentinel
cfg := config.NewDefaultConfig()
cfg.Sentinel.CircuitBreakerStrategy = circuitbreaker.StrategyErrorCountRatio // 或 ErrorRatio
config.ApplyConfig(cfg)

// 熔断规则:基于错误比率
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,            // 慢调用阈值
}
circuitbreaker.PutRule(rule)

Sentinel 支持的三种熔断策略

策略说明适用场景
ErrorCount单位时间内错误数超过阈值错误绝对数量可控的场景
ErrorRatio错误率超过阈值按比例熔断,更灵活
SlowRatio慢调用比例超过阈值应对响应变慢但不报错的情况

在 Kitex 中使用

import "github.com/alibaba/sentinel-golang/adapters/kitex"

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 秒,平滑统计 ──────→
type SlidingWindowCounter struct {
	interval     time.Duration
	buckets      []atomic.Int64
	bucketCount  int
 currentIndex int
}

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,
	}
}

func (sw *SlidingWindowCounter) currentBucketIndex() int {
	return int(time.Now().UnixNano() / sw.interval.Nanoseconds()) % sw.bucketCount
}

func (sw *SlidingWindowCounter) Increment(isSuccess bool) {
	idx := sw.currentBucketIndex()
	if isSuccess {
		sw.buckets[idx].Add(1)
	} else {
		sw.buckets[idx].Add(-1) // 负数表示失败
	}
}

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)
}

熔断 vs 限流 vs 超时 vs 重试

这四者经常配合使用,分工明确:

机制层级目的触发条件行为
超时客户端防止等待过久超过指定时间未返回放弃等待,返回错误
限流服务端保护自身不被打垮超过 QPS/并发上限拒绝新请求
熔断客户端防止下游故障扩散错误率/慢调用率过高切断对下游的调用
重试客户端提高成功率临时性失败(网络抖动等)再次发起请求
请求流程:
  客户端
    ├─ [1] 熔断检查 ── 熔断打开 ──→ 快速失败(不浪费资源)
    │        ↓ 允许
    ├─ [2] 超时控制 ── 设定 deadline
    │        ↓
    ├─ [3] 负载均衡 ── 选一个服务端实例
    │        ↓
    ├─ [4] 发送请求 ── 等待响应
    │        ↓
    ├─ [5] 收到响应 ── 成功?
    │     ╱          ╲
    │   成功         失败
    │    │           │
    │    ↓      ┌────┴─────┐
    │         [6] 重试?    │
    │         有剩余次数?   │
    │       ╱      ╲       │
    │     是       否      │
    │     ↓        ↓       │
    │   回到[4]  返回错误  │
    │                    ↓
    │                 [7] 熔断器记录失败

最佳实践

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)

熔断不是终点,配合降级逻辑提供更友好的用户体验:

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)
			}

			err = next(ctx, req, resp)
			if err != nil {
				cb.RecordFailure()
			} else {
				cb.RecordSuccess()
			}
			return err
		}
	}
}

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)
}

注意事项

  1. 不要对幂等性差的操作过度重试:熔断 + 重试组合时,注意重试可能放大流量
  2. 熔断是客户端行为:每个客户端独立维护熔断状态,可能出现部分客户端已恢复而另一部分仍在熔断的情况
  3. 半开不是万能的:如果下游真的挂了,半开试探只会增加无效请求,设置合理的 PassDuration 很关键
  4. 区分错误类型:不是所有错误都应该计入熔断(如参数校验错误不应触发熔断),建议在 RecordFailure 前过滤
  5. 连接池影响:熔断器关闭后,如果连接池也断了,第一个请求仍然可能超时,建议配合 [[预热]] 使用

相关笔记

  • [[Kitex/限流]] — 服务端限流,与熔断互补
  • [[Kitex/超时]] — 熔断的前提配置,没有超时熔断无从谈起
  • [[Kitex/重试]] — 熔断 + 重试的组合策略
  • [[Kitex/负载均衡]] — 负载均衡 + 熔断可实现自动摘除故障节点
  • [[Kitex/中间件]] — 熔断通过 Middleware 机制嵌入
About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!