Appearance
熔断、降级与限流
本篇是 Go 微服务系列的第八篇。微服务系统中故障是常态:依赖的下游服务可能超时、宕机、变慢。如果不加保护,一个慢服务会拖垮整条调用链,造成「雪崩」。稳定性保障的核心三件套——熔断、降级、限流——就是应对这类问题的工具。本篇将讲解原理、用 sony/gobreaker、golang.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 解耦微服务、削峰填谷。
延伸阅读:
- gobreaker:https://github.com/sony/gobreaker
- go-kit circuitbreaker:https://pkg.go.dev/github.com/go-kit/kit/circuitbreaker
- Release It! 第二版,Michael Nygard 著
- Google SRE Book:https://sre.google/sre-book/