Skip to content

熔断、降级与限流

本篇是 Go 微服务系列的第八篇。微服务系统中故障是常态:依赖的下游服务可能超时、宕机、变慢。如果不加保护,一个慢服务会拖垮整条调用链,造成「雪崩」。稳定性保障的核心三件套——熔断、降级、限流——就是应对这类问题的工具。本篇将讲解原理、用 sony/gobreakergolang.org/x/time/rate、go-kit 的断路器,并给出完整可运行示例。

一、稳定性问题的本质

1. 雪崩效应

考虑一个典型场景:订单服务调用库存服务,库存服务变慢(响应从 50ms 涨到 2s)。

  • 订单服务调用库存的请求堆积,goroutine 阻塞。
  • 订单服务自己的 goroutine 池被打满,对新请求也无法响应。
  • 上游网关调用订单服务的请求堆积。
  • 整个链路雪崩。

2. 三种保护手段

  • 超时控制:每个调用都有超时,避免无限阻塞。
  • 熔断器:当下游错误率超阈值,自动「断开」一段时间,快速失败。
  • 限流器:限制每秒处理的请求数,防止被打垮。
  • 降级:核心链路出问题时,提供兜底响应(默认值、缓存、空响应)。

四者通常组合使用:

请求 ──► 限流(防过载) ──► 熔断(防级联故障) ──► 超时(防阻塞) ──► 下游
                                                       │ 失败

                                                    降级返回

二、熔断器模式原理

1. 熔断器的三个状态

熔断器(Circuit Breaker)有三个状态:

  • Closed(关闭):正常放行请求,统计错误率。
  • Open(打开):错误率达到阈值,直接快速失败,不发请求。
  • HalfOpen(半开):经过冷却期后,放行少量试探请求,根据结果决定回到 Closed 还是 Open。

状态机:

        error rate > threshold              cooldown elapsed
       ┌──────────────────────────┐      ┌──────────────────────┐
       ▼                          │      ▼                      │
   ┌────────┐  trial success   ┌────────┐  trial fail   ┌──────┴───┐
   │  Open  │ ◄─────────────── │HalfOpen│ ────────────► │  Open    │
   └────────┘                  └────────┘                └──────────┘
       ▲                          │
       │ trial fail                │ trial success
       │                          ▼
       └─────────────────────── ┌────────┐
                                │ Closed │
                                └────────┘
                                  ▲ │
                                  │ │ error rate > threshold
                                  └─┘ (重新计数)

2. 关键参数

  • 错误率阈值:触发 Open 的错误率(如 50%)。
  • 最小请求数:达到这个数才计算错误率(避免样本过小误判)。
  • 冷却时间:Open 状态持续时间(如 30s)。
  • 半开试探数:HalfOpen 放行的请求数(如 5)。
  • 成功阈值:HalfOpen 阶段连续成功多少次才回到 Closed。

3. 熔断 vs 重试

熔断器打开时不要重试——重试只会雪上加霜。重试只在 Closed 状态、偶发错误时使用。

三、go-kit 的断路器

go-kit 的 circuitbreaker 子包提供了对多个断路器实现的统一封装(包括 gobreaker、Hystrix、handy):

go
package main

import (
	"context"
	"errors"
	"fmt"
	"log"
	"time"

	"github.com/go-kit/kit/circuitbreaker"
	"github.com/sony/gobreaker"
)

type Service interface {
	Call(ctx context.Context, s string) (string, error)
}

type myService struct{}

func (s *myService) Call(_ context.Context, in string) (string, error) {
	if in == "fail" {
		return "", errors.New("upstream error")
	}
	return "ok: " + in, nil
}

// 用 gobreaker 包裹 Service
func withCircuitBreaker(svc Service) Service {
	cb := gobreaker.NewCircuitBreaker(gobreaker.Settings{
		Name:        "my-svc",
		MaxRequests: 1,
		Interval:    10 * time.Second,
		Timeout:     5 * time.Second,
		ReadyToTrip: func(counts gobreaker.Counts) bool {
			failureRatio := float64(counts.TotalFailures) / float64(counts.Requests)
			return counts.Requests > 5 && failureRatio > 0.5
		},
		OnStateChange: func(name string, from, to gobreaker.State) {
			log.Printf("[cb] %s: %s -> %s", name, from, to)
		},
	})
	return &cbService{cb: cb, inner: svc}
}

type cbService struct {
	cb    *gobreaker.CircuitBreaker
	inner Service
}

func (s *cbService) Call(ctx context.Context, in string) (string, error) {
	v, err := s.cb.Execute(func() (interface{}, error) {
		return s.inner.Call(ctx, in)
	})
	if err != nil {
		return "", err
	}
	return v.(string), nil
}

func main() {
	svc := withCircuitBreaker(&myService{})
	for i := 0; i < 20; i++ {
		_, err := svc.Call(context.Background(), "fail")
		fmt.Printf("call %d: err=%v\n", i, err)
		time.Sleep(100 * time.Millisecond)
	}
	// 熔断后即使调用 ok 也会快速失败
	_, err := svc.Call(context.Background(), "ok")
	fmt.Printf("after open, call ok: %v\n", err)
	_ = circuitbreaker.Gobreaker // 引用以保留 import
}

go-kit 的好处是把断路器抽象成 middleware,可以和它的 endpoint / transport 模式无缝组合。

四、sony/gobreaker 实战

sony/gobreaker 是 Go 生态最流行的断路器实现之一,API 简洁,功能完备。

1. 基础用法

go
package main

import (
	"context"
	"errors"
	"fmt"
	"log"
	"math/rand"
	"time"

	"github.com/sony/gobreaker"
)

type UserService struct {
	failureRate float64
}

func (s *UserService) GetUser(ctx context.Context, id int) (string, error) {
	if rand.Float64() < s.failureRate {
		return "", errors.New("upstream error")
	}
	return fmt.Sprintf("user-%d", id), nil
}

func newBreaker(name string) *gobreaker.CircuitBreaker {
	return gobreaker.NewCircuitBreaker(gobreaker.Settings{
		Name:        name,
		MaxRequests: 3,                  // 半开状态最多放行 3 个
		Interval:    60 * time.Second,   // 计数窗口
		Timeout:     5 * time.Second,    // open 持续时间
		ReadyToTrip: func(c gobreaker.Counts) bool {
			// 至少 10 个请求且错误率 > 60% 才熔断
			failureRatio := float64(c.TotalFailures) / float64(c.Requests)
			return c.Requests >= 10 && failureRatio > 0.6
		},
		OnStateChange: func(name string, from, to gobreaker.State) {
			log.Printf("[breaker:%s] %s -> %s", name, from, to)
		},
		IsSuccessful: func(err error) bool {
			// 可以把某些业务错误视为成功(不计数为失败)
			return err == nil
		},
	})
}

func main() {
	rand.Seed(time.Now().UnixNano())
	svc := &UserService{failureRate: 0.7}
	cb := newBreaker("user-service")

	for i := 0; i < 30; i++ {
		result, err := cb.Execute(func() (interface{}, error) {
			return svc.GetUser(context.Background(), i)
		})
		state := cb.State()
		log.Printf("call %2d  state=%-9s  result=%v  err=%v", i, state, result, err)
		time.Sleep(200 * time.Millisecond)
	}
}

运行后你会看到状态从 Closed 变为 Open,过 5 秒后变为 HalfOpen 试探,再根据结果切换。

2. 区分错误类型

实际业务中,业务错误(如 NotFound)不应触发熔断,只有系统错误(如超时、连接失败)才应触发。用 IsSuccessful 自定义:

go
cb := gobreaker.NewCircuitBreaker(gobreaker.Settings{
    Name: "user-service",
    ReadyToTrip: func(c gobreaker.Counts) bool {
        return float64(c.TotalFailures)/float64(c.Requests) > 0.5
    },
    IsSuccessful: func(err error) bool {
        // 业务错误不计入熔断
        var bizErr *BizError
        if errors.As(err, &bizErr) {
            return true
        }
        return err == nil
    },
})

五、限流算法

1. 令牌桶(Token Bucket)

  • 桶以恒定速率生成令牌(容量上限)。
  • 请求到达时取一个令牌:有则放行,无则拒绝。
  • 允许一定突发流量(桶满时可一次性放行多个)。

适合:允许突发、平滑限流的场景。Go 标准库 golang.org/x/time/rate 实现的就是令牌桶

2. 漏桶(Leaky Bucket)

  • 请求如水滴进桶,桶以恒定速率漏出(处理)。
  • 桶满则拒绝。
  • 输出速率严格恒定。

适合:要求严格平滑的场景。

3. 滑动窗口

  • 把时间分成小窗口,统计当前窗口内请求数。
  • 滑动合并多个小窗口,平滑计数。

适合:精确控制每秒请求数。常见于网关的 QPS 限流。

4. 算法对比

算法突发流量实现复杂度平滑度典型场景
令牌桶允许通用
漏桶不允许极高整流
滑动窗口限制精确 QPS
计数器不允许简单限流(不推荐)

六、golang.org/x/time/rate 限流

1. 基础用法

go
package main

import (
	"context"
	"fmt"
	"time"

	"golang.org/x/time/rate"
)

func main() {
	// rate=5 每秒 5 个令牌,burst=10 桶容量
	limiter := rate.NewLimiter(5, 10)

	for i := 0; i < 20; i++ {
		// Allow:拿不到立即返回 false(非阻塞)
		if !limiter.Allow() {
			fmt.Printf("req %d: rejected\n", i)
			continue
		}
		fmt.Printf("req %d: allowed\n", i)
	}

	// Wait:拿不到就阻塞等待(带 ctx 超时)
	ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
	defer cancel()
	for i := 0; i < 20; i++ {
		if err := limiter.Wait(ctx); err != nil {
			fmt.Printf("req %d: wait failed: %v\n", i, err)
			return
		}
		fmt.Printf("req %d: passed after wait\n", i)
	}
}

三个核心 API:

  • Allow() / AllowN(now, n):非阻塞,立即返回是否放行。
  • Wait(ctx) / WaitN(ctx, n):阻塞等待令牌,支持 context 超时。
  • Reserve() / ReserveN(now, n):预订令牌,返回需要等待的时间。

2. 多维度限流

实际业务常需要按用户、按 IP、按 API 维度限流。维护一个 limiter 池:

go
package main

import (
	"sync"
	"time"

	"golang.org/x/time/rate"
)

// 注意:实际导入路径是 golang.org/x/time/rate,下面为示例
type RateLimiter struct {
	mu       sync.Mutex
	limiters map[string]*entry
	rate     rate.Limit
	burst    int
	ttl      time.Duration
}

type entry struct {
	limiter   *rate.Limiter
	lastSeen  time.Time
}

func NewRateLimiter(r rate.Limit, b int) *RateLimiter {
	rl := &RateLimiter{
		limiters: make(map[string]*entry),
		rate:     r,
		burst:    b,
		ttl:      10 * time.Minute,
	}
	go rl.gc()
	return rl
}

func (rl *RateLimiter) Get(key string) *rate.Limiter {
	rl.mu.Lock()
	defer rl.mu.Unlock()
	if e, ok := rl.limiters[key]; ok {
		e.lastSeen = time.Now()
		return e.limiter
	}
	l := rate.NewLimiter(rl.rate, rl.burst)
	rl.limiters[key] = &entry{limiter: l, lastSeen: time.Now()}
	return l
}

func (rl *RateLimiter) gc() {
	for {
		time.Sleep(time.Minute)
		rl.mu.Lock()
		for k, e := range rl.limiters {
			if time.Since(e.lastSeen) > rl.ttl {
				delete(rl.limiters, k)
			}
		}
		rl.mu.Unlock()
	}
}

func main() {
	// 仅供编译演示:实际用 import "golang.org/x/time/rate"
	// rl := NewRateLimiter(rate.Limit(10), 20)
	// _ = rl.Get("user-123").Allow()
}

注意上面为了能独立编译演示,省略了真实导入。实际项目里直接用 golang.org/x/time/rate 即可。

七、超时控制策略

1. context 超时

Go 的 context 是超时控制的基础:

go
package main

import (
	"context"
	"fmt"
	"time"
)

func callDownstream(ctx context.Context) (string, error) {
	// 模拟下游慢
	select {
	case <-time.After(2 * time.Second):
		return "ok", nil
	case <-ctx.Done():
		return "", ctx.Err()
	}
}

func handler() {
	// 总超时 1 秒
	ctx, cancel := context.WithTimeout(context.Background(), 1*time.Second)
	defer cancel()

	result, err := callDownstream(ctx)
	fmt.Printf("result=%q err=%v\n", result, err)
}

func main() {
	handler()
}

2. 多级超时预算

一个请求可能调用多个下游,要给每个分配预算。常见做法:

  • 入口请求设总超时(如 5s)。
  • 把剩余预算传递给下游调用。
  • 下游调用至少保留 100ms 给自己处理。
go
package main

import (
	"context"
	"fmt"
	"time"
)

func callService(ctx context.Context, name string, duration time.Duration) (string, error) {
	deadline, ok := ctx.Deadline()
	if ok {
		remaining := time.Until(deadline)
		fmt.Printf("[%s] remaining budget: %v\n", name, remaining)
		if remaining < duration+200*time.Millisecond {
			return "", fmt.Errorf("%s: insufficient budget", name)
		}
	}
	select {
	case <-time.After(duration):
		return name + ":ok", nil
	case <-ctx.Done():
		return "", ctx.Err()
	}
}

func main() {
	// 总预算 3 秒
	ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
	defer cancel()

	r1, err := callService(ctx, "user", 1*time.Second)
	fmt.Printf("user: %v %v\n", r1, err)

	r2, err := callService(ctx, "order", 1*time.Second)
	fmt.Printf("order: %v %v\n", r2, err)

	r3, err := callService(ctx, "pay", 2*time.Second)
	fmt.Printf("pay: %v %v\n", r3, err)
}

3. 超时设置的坑

  • 超时过短:正常请求被切断,重试洪流。
  • 超时过长:失去保护意义。
  • 超时层级错配:上游超时短于下游,下游还在跑上游就放弃了,浪费资源。

经验值:上游超时 = 下游超时 × 1.5,给下游处理和重试留余地。

八、降级策略设计

1. 什么是降级

降级(Degradation)指系统压力大或部分功能不可用时,主动牺牲非核心功能,保证核心功能可用。

2. 常见降级策略

  • 返回默认值:推荐列表服务挂了,返回热门商品。
  • 返回缓存:数据库挂了,返回缓存数据(可能稍旧)。
  • 返回部分结果:聚合接口某个下游失败,返回部分字段。
  • 同步转异步:支付场景写消息队列异步处理。
  • 关闭非核心功能:双 11 关闭评论、推荐。

3. 降级示例

go
package main

import (
	"context"
	"errors"
	"fmt"
	"time"
)

type Recommendation struct {
	Items []string
}

// RecommendFromRemote 调用远程推荐服务
func RecommendFromRemote(ctx context.Context, userID string) (*Recommendation, error) {
	// 模拟偶尔失败
	if time.Now().Unix()%2 == 0 {
		return nil, errors.New("remote unavailable")
	}
	return &Recommendation{Items: []string{"remote-a", "remote-b"}}, nil
}

// RecommendFromCache 从缓存读
func RecommendFromCache(ctx context.Context, userID string) (*Recommendation, error) {
	return &Recommendation{Items: []string{"cached-a", "cached-b"}}, nil
}

// RecommendDefault 默认兜底
func RecommendDefault() *Recommendation {
	return &Recommendation{Items: []string{"default-1", "default-2", "default-3"}}
}

// GetRecommendation 多级降级:远程 → 缓存 → 默认
func GetRecommendation(ctx context.Context, userID string) *Recommendation {
	ctx, cancel := context.WithTimeout(ctx, 200*time.Millisecond)
	defer cancel()

	if r, err := RecommendFromRemote(ctx, userID); err == nil {
		fmt.Println("from remote")
		return r
	}
	if r, err := RecommendFromCache(ctx, userID); err == nil {
		fmt.Println("from cache")
		return r
	}
	fmt.Println("from default")
	return RecommendDefault()
}

func main() {
	r := GetRecommendation(context.Background(), "user-1")
	fmt.Printf("items: %v\n", r.Items)
}

4. 降级的开关

降级不应写死在代码里,应通过配置中心控制:

  • 自动降级:熔断器打开时自动降级。
  • 手动降级:运维通过配置中心下发开关。

参考第 6 篇的配置中心实现,把降级开关做成配置项即可。

九、完整示例:带熔断和限流的服务调用

下面给出一个完整的可运行示例,集成熔断器、限流器、超时控制、降级:

go
package main

import (
	"context"
	"errors"
	"fmt"
	"log"
	"math/rand"
	"sync/atomic"
	"time"

	"github.com/sony/gobreaker"
	"golang.org/x/time/rate"
)

// === 下游服务(模拟一个不稳定的服务)===

type Downstream struct {
	failureRate float64
	maxLatency  time.Duration
}

func (d *Downstream) Call(ctx context.Context, req string) (string, error) {
	latency := time.Duration(rand.Intn(int(d.maxLatency)))
	select {
	case <-time.After(latency):
		if rand.Float64() < d.failureRate {
			return "", errors.New("upstream error")
		}
		return "ok: " + req, nil
	case <-ctx.Done():
		return "", ctx.Err()
	}
}

// === 熔断器封装 ===

type ProtectedClient struct {
	cb       *gobreaker.CircuitBreaker
	limiter  *rate.Limiter
	timeout  time.Duration
	fallback func() (string, error)
}

func NewProtectedClient(name string, rps rate.Limit, burst int, timeout time.Duration) *ProtectedClient {
	cb := gobreaker.NewCircuitBreaker(gobreaker.Settings{
		Name:        name,
		MaxRequests: 3,
		Interval:    30 * time.Second,
		Timeout:     5 * time.Second,
		ReadyToTrip: func(c gobreaker.Counts) bool {
			failureRatio := float64(c.TotalFailures) / float64(c.Requests)
			return c.Requests >= 5 && failureRatio > 0.5
		},
		OnStateChange: func(name string, from, to gobreaker.State) {
			log.Printf("[cb:%s] %s -> %s", name, from, to)
		},
	})
	return &ProtectedClient{
		cb:       cb,
		limiter:  rate.NewLimiter(rps, burst),
		timeout:  timeout,
		fallback: func() (string, error) {
			return "fallback: default-value", nil
		},
	}
}

func (p *ProtectedClient) Call(ctx context.Context, ds *Downstream, req string) (string, error) {
	// 1. 限流
	if !p.limiter.Allow() {
		atomic.AddInt64(&metricRateLimited, 1)
		return "", errors.New("rate limited")
	}

	// 2. 超时预算
	ctx, cancel := context.WithTimeout(ctx, p.timeout)
	defer cancel()

	// 3. 熔断
	result, err := p.cb.Execute(func() (interface{}, error) {
		return ds.Call(ctx, req)
	})
	if err != nil {
		// 熔断打开或调用失败 → 降级
		atomic.AddInt64(&metricFallback, 1)
		return p.fallback()
	}
	atomic.AddInt64(&metricSuccess, 1)
	return result.(string), nil
}

// === 指标 ===

var (
	metricSuccess    int64
	metricFallback   int64
	metricRateLimited int64
)

func printMetrics() {
	log.Printf("[metrics] success=%d fallback=%d rate_limited=%d",
		atomic.LoadInt64(&metricSuccess),
		atomic.LoadInt64(&metricFallback),
		atomic.LoadInt64(&metricRateLimited))
}

func main() {
	rand.Seed(time.Now().UnixNano())
	ds := &Downstream{failureRate: 0.6, maxLatency: 300 * time.Millisecond}
	client := NewProtectedClient("order-service", rate.Limit(20), 5, 500*time.Millisecond)

	// 模拟 50 个并发请求
	var wg sync.WaitGroup
	for i := 0; i < 50; i++ {
		wg.Add(1)
		go func(idx int) {
			defer wg.Done()
			time.Sleep(time.Duration(rand.Intn(500)) * time.Millisecond)
			result, err := client.Call(context.Background(), ds, fmt.Sprintf("req-%d", idx))
			log.Printf("req %2d: result=%s err=%v", idx, result, err)
		}(i)
	}
	wg.Wait()
	printMetrics()
}

需要补一个 sync 的 import,完整 import 块为:

go
import (
    "context"
    "errors"
    "fmt"
    "log"
    "math/rand"
    "sync"
    "sync/atomic"
    "time"

    "github.com/sony/gobreaker"
    "golang.org/x/time/rate"
)

运行这个例子,你会观察到:

  • 一部分请求成功(success 计数)。
  • 一部分被限流(rate_limited 计数)。
  • 一部分触发熔断并降级(fallback 计数)。
  • 日志中能看到熔断器状态切换:Closed → Open → HalfOpen → Closed。

这就是一个生产可用的稳定性保护骨架。

十、小结

本篇我们学习了微服务稳定性保障的三件套:

  • 熔断器:错误率超阈值后快速失败,防止级联故障;三状态机 Closed → Open → HalfOpen。
  • 限流:令牌桶(x/time/rate)允许突发,漏桶严格整流,滑动窗口精确 QPS。
  • 超时控制:context 传递超时,多级调用要做预算分配。
  • 降级:远程失败 → 缓存 → 默认值,逐级兜底,保核心可用。
  • gobreaker:Go 生态最流行的断路器,支持 IsSuccessful 区分业务错误。
  • go-kit circuitbreaker:统一封装多个断路器实现,与 endpoint 模式整合。
  • 组合使用:限流在最外层,熔断在中层,超时在内层,降级兜底,构成完整的稳定性防线。

下一篇我们将进入消息队列与异步通信,学习用 NATS 和 Kafka 解耦微服务、削峰填谷。

延伸阅读