学习目标
类比(先建立直觉):微服务的"可用性"和"可观测性"就像你开了一家连锁餐厅。可用性是"出事时还能接住客人"——后厨着火就挂"暂停营业"(熔断)、客人太多就取号排队(限流)、人手不够就只卖预制菜(降级);可观测性是"装监控、对讲机、点单日志"——出了差错能立刻定位到是哪桌、哪道菜、哪个厨师。本章就是教你怎么把这套"餐厅应急 + 监控"机制用代码落到微服务框架里。
学完本章你应该能够:
- 用自己的话讲清熔断、限流、降级三者"本质相同、处理策略不同"的关系,并画出故障处理流程与熔断器状态机。
- 讲清令牌桶、漏桶、固定窗口、滑动窗口四种限流算法的原理、差异(尤其是临界点问题、突发流量)与适用场景。
- 基于 gRPC Interceptor 写一个可用的限流 / 熔断拦截器,并说清客户端限流 vs 服务端限流、单机限流 vs 集群限流的区别。
- 讲清可观测性三支柱(Metrics / Tracing / Logging)的分工,以及 TraceID / SpanID 如何把一次跨服务调用串联成完整链路。
- 说出"基于可观测性的服务治理“的基本思路(指标采集 → 客户端拉取 / 平台推送 → 驱动负载均衡与限流),并在面试中讲成一个完整故事。
前置知识:
- 第 1–4 章:网络编程、RPC 协议、服务注册发现、负载均衡与集群容错
- Go 基础:
goroutine、sync.Mutex、context、time - gRPC 基本用法与 Interceptor 机制
- Redis 基础命令:
INCR/ZADD/ Lua 脚本
本章你会动手做的事:
- 逐一跑一遍令牌桶 / 漏桶 / 固定窗口 / 滑动窗口的
Allow(),观察高并发下是否如预期拒绝请求。 - 写一个 gRPC 服务端限流拦截器(令牌桶),对同一方法做 QPS 限制。
- 本地起 Redis,用 Lua 脚本实现一次集群固定窗口限流,压测验证其原子性。
前言:微服务的稳定运行与问题排查
在前面的章节中,我们已经构建了一个完整的微服务框架——从网络编程到 RPC 协议,从服务注册发现到负载均衡和集群容错。但是,一个微服务框架光能跑起来还不够,还需要解决两个关键问题:
- 可用性——当服务出现故障或流量激增时,如何保证系统仍然可用?
- 可观测性——当系统出问题时,如何快速定位和排查?
本教程涵盖以下主题:
- AOP 方案对比——Kratos、Dubbo-go、go-micro 的拦截器设计
- 可用性核心概念——熔断、限流、降级的关系与区别
- 限流算法详解——令牌桶、漏桶、固定窗口、滑动窗口及 Go 实现
- 熔断与降级——gRPC 拦截器实现限流与熔断
- 集群限流——基于 Redis 的分布式限流
- 可观测性三支柱——Metrics(指标)、Tracing(链路追踪)、Logging(日志)
- 基于可观测性的服务治理——如何利用观测数据驱动治理决策
- 面试要点总结
一、AOP 方案对比
在微服务框架中,可用性和可观测性的功能都是通过 AOP(面向切面编程) 来实现的——即在请求处理的前后插入额外的逻辑,而不侵入业务代码。不同框架的 AOP 方案各有特色。
1.1 Kratos 的 AOP 方案
Kratos 使用的是类似洋葱模型的中间件设计,我们在 Web 框架和 ORM 框架中都见过类似的设计。核心是 Middleware 接口,请求像穿过洋葱一样经过每一层中间件。
1.2 Dubbo-go 的 AOP 方案
Dubbo-go 的设计更加接近责任链模式。其中它有一个回调接口 OnResponse,即响应返回的时候会执行这个方法。
// Dubbo-go 的 Filter 接口(概念示例)
// 请求和响应分别走两个方向,形成完整的调用链
type Filter interface {
// Invoke 处理请求(正向调用)
Invoke(invoker Invoker, invocation Invocation) Result
// OnResponse 处理响应(反向回调)
OnResponse(result Result, invoker Invoker, invocation Invocation) Result
}
1.3 go-micro 的 AOP 方案
go-micro 叫做 Wrapper,也就是认为自己是在已有功能的基础上再封装一些功能,设计者认为它更加接近装饰器模式或者洋葱模式。
客户端 Wrapper 分成三类:
| 类型 | 说明 |
|---|---|
CallWrapper | 最细粒度控制,针对每一次调用 |
Wrapper | 针对客户端的 Wrapper |
StreamWrapper | 针对 Stream API 的 Wrapper |
实际上
Wrapper和CallWrapper只需要保留一个就可以。通过注册过程来控制作用范围,例如我们在某个调用里面注册 Wrapper,那么只对该调用起效果。
服务端的 Wrapper 定义不太一样,核心就是 HandlerWrapper 和 StreamWrapper,其作用接近于客户端的 Wrapper 和 StreamWrapper。
设计偏好:个人比较喜欢 Dubbo-go 那种一致的抽象——客户端和服务端使用相同的 API,普通请求和 Stream 请求统一对待。
1.4 gRPC 的 Interceptor
gRPC 的拦截器叫做 Interceptor,分成好几种:
| 类型 | 说明 |
|---|---|
UnaryServerInterceptor | 服务端一元请求拦截器 |
StreamServerInterceptor | 服务端流请求拦截器 |
UnaryClientInterceptor | 客户端一元请求拦截器 |
StreamClientInterceptor | 客户端流请求拦截器 |
拦截器的原理和我们前面使用的 Middleware 设计是一样的,只是叫法不同。我们后续的限流、熔断、可观测性代码都基于 gRPC 的 Interceptor 来实现。
二、可用性核心概念
2.1 可用性的五大主题
在微服务框架中,可用性是和服务治理最密切相关的主题:
| 主题 | 说明 | 前序章节 |
|---|---|---|
| 熔断 | 当服务故障时,拒绝新的请求,防止故障蔓延 | 本章详解 |
| 限流 | 一段时间内只允许特定数量的请求被处理 | 本章详解 |
| 降级 | 全部请求执行一段更加简单的逻辑(走快路径) | 本章详解 |
| 重试 | 调用失败后重试,可能换节点 | 已在 Cluster 章节讨论 |
| 超时控制 | 设置请求超时时间,防止无限等待 | 已在 RPC 协议章节讨论 |
2.2 熔断、限流、降级的关系
熔断、限流和降级并没有本质区别,都可以归属到故障处理的范畴里面:
flowchart TB
A[故障检测
判定服务是否不健康] --> B{故障处理}
B -->|只允许一部分请求通过| C[限流]
B -->|全部请求被拒绝| D[熔断]
B -->|全部请求走简易逻辑| E[降级]
C --> F[故障恢复
过一段时间恢复]
D --> F
E --> F故障处理的三个层次:
- 限流:一段时间内只允许特定数量的请求被正常处理
- 熔断:全部请求都会被拒绝(返回错误)
- 降级:全部请求都会执行一段更加简单的逻辑(返回默认值或走快路径)
三者的区别其实很少:
- 如果判断到资源不足,只允许一部分请求被正常处理,那就是限流
- 对于没有被正常处理的请求,如果直接拒绝返回错误,那就是被熔断了
- 如果请求没有被拒绝,而是返回了默认值或走了简易路径,那就是被降级了
- 从接口设计的角度来说,它们也非常接近,例如定义一个
Allow方法就可以了- 算法层面上它们也都是通用的
2.3 故障检测算法分类
从故障检测算法的类型来看,可以分成两类:
| 类型 | 说明 | 算法示例 |
|---|---|---|
| 静态类型 | 依赖于测试或程序员经验提前设置阈值 | 令牌桶、漏桶、固定窗口、滑动窗口 |
| 动态类型 | 根据服务当前状态动态判断 | 基于错误率、响应时间、BBR 算法 |
对于绝大多数应用来说,静态类型算法就足够了。
三、限流算法详解
3.1 令牌桶(Token Bucket)
原理:有一个"人"按一定的速率发令牌,令牌会被放到一个桶里。每一个请求从桶里面拿一个令牌,拿到令牌的请求就会被处理,没有拿到令牌的请求就会被拒绝或阻塞。
flowchart LR
G[令牌生成器
按固定速率生成] -->|放入令牌| B[令牌桶
容量有限]
R1[请求1] -->|拿令牌| B
R2[请求2] -->|拿令牌| B
R3[请求3] -->|令牌不够| X[拒绝/阻塞]
B -->|有令牌| R1
B -->|有令牌| R2要点:
- 有一个人按一定的速率发令牌
- 令牌会被放到一个桶里(桶有容量上限)
- 每一个请求从桶里面拿一个令牌
- 拿到令牌的请求就会被处理
- 没有拿到令牌的请求就会:直接被拒绝,或者阻塞直到拿到令牌或者超时
特点:令牌桶允许一定程度的突发流量——如果桶里积累了很多令牌,突然来一波请求可以一次性处理多个。
package ratelimit
import (
"sync"
"time"
)
// ============================================================
// 令牌桶(Token Bucket)限流算法实现
// ============================================================
// TokenBucket 令牌桶限流器
type TokenBucket struct {
mu sync.Mutex // 互斥锁,保证并发安全
rate float64 // 令牌生成速率(个/秒)
capacity float64 // 桶的容量(最大令牌数)
tokens float64 // 当前桶中的令牌数
lastRefill time.Time // 上次补充令牌的时间
}
// NewTokenBucket 创建一个令牌桶限流器
// 参数 rate 是令牌生成速率(每秒生成多少个令牌)
// 参数 capacity 是桶的容量(最多存放多少个令牌)
func NewTokenBucket(rate, capacity float64) *TokenBucket {
return &TokenBucket{
rate: rate, // 例如 10.0 表示每秒生成 10 个令牌
capacity: capacity, // 例如 100.0 表示桶最多放 100 个令牌
tokens: capacity, // 初始化时桶是满的
lastRefill: time.Now(), // 记录创建时间
}
}
// Allow 尝试获取一个令牌
// 返回 true 表示获取成功(请求被允许),false 表示获取失败(请求被拒绝)
func (tb *TokenBucket) Allow() bool {
tb.mu.Lock()
defer tb.mu.Unlock()
// 第一步:计算自上次补充以来经过的时间,补充令牌
now := time.Now()
elapsed := now.Sub(tb.lastRefill).Seconds() // 经过的秒数
tb.lastRefill = now
// 按速率补充令牌,但不超过桶的容量
// 例如 rate=10,elapsed=0.5秒,则补充 10*0.5=5 个令牌
tb.tokens += tb.rate * elapsed
if tb.tokens > tb.capacity {
tb.tokens = tb.capacity // 令牌数不能超过桶容量
}
// 第二步:尝试消费一个令牌
if tb.tokens >= 1 {
tb.tokens-- // 消费一个令牌
return true
}
// 令牌不足,拒绝请求
return false
}
3.2 漏桶(Leaky Bucket)
原理:请求过来先排队,每隔一段时间放过去一个请求,请求排队直到通过或者超时。
flowchart LR
R1[请求1] --> Q[请求队列
漏桶]
R2[请求2] --> Q
R3[请求3] --> Q
Q -->|每隔一段时间放一个| O1[处理请求1]
Q -->|队列满| X[拒绝请求3]要点:
- 请求过来先排队
- 每隔一段时间,放过去一个请求
- 请求排队直到通过,或者超时
漏桶与令牌桶效果是一样的。令牌桶是"请求拿令牌才能通过”,漏桶是"请求排队等待处理"。两者都可以实现固定速率的流量控制。
特点:漏桶的输出速率是严格固定的,不管来了多少请求,处理速度始终一样。这适合需要严格控速的场景。
package ratelimit
import (
"sync"
"time"
)
// ============================================================
// 漏桶(Leaky Bucket)限流算法实现
// ============================================================
// LeakyBucket 漏桶限流器
type LeakyBucket struct {
mu sync.Mutex
rate float64 // 漏水速率(请求/秒)
capacity float64 // 桶容量(队列长度上限)
water float64 // 当前桶中的水量(排队请求数)
lastLeak time.Time // 上次漏水时间
}
// NewLeakyBucket 创建一个漏桶限流器
// 参数 rate 是处理速率(每秒处理多少请求)
// 参数 capacity 是桶容量(最多排队多少请求)
func NewLeakyBucket(rate, capacity float64) *LeakyBucket {
return &LeakyBucket{
rate: rate, // 例如 10.0 表示每秒处理 10 个请求
capacity: capacity, // 例如 100.0 表示最多排队 100 个请求
water: 0, // 初始水量为 0
lastLeak: time.Now(),
}
}
// Allow 尝试将请求加入队列
// 返回 true 表示请求被接受(排队等待处理),false 表示队列已满(拒绝)
func (lb *LeakyBucket) Allow() bool {
lb.mu.Lock()
defer lb.mu.Unlock()
// 第一步:先漏水(处理排队请求)
now := time.Now()
elapsed := now.Sub(lb.lastLeak).Seconds()
lb.lastLeak = now
// 按速率漏水,水量不能小于 0
// 例如 rate=10,elapsed=0.5秒,则漏掉 10*0.5=5 个请求
lb.water -= lb.rate * elapsed
if lb.water < 0 {
lb.water = 0
}
// 第二步:尝试加入新请求(加水)
if lb.water < lb.capacity {
lb.water++ // 加入队列
return true
}
// 队列已满,拒绝请求
return false
}
3.3 固定窗口(Fixed Window)
原理:将时间划分为固定大小的窗口(如每分钟一个窗口),在每个窗口内统计请求数,超过阈值就拒绝。
flowchart LR
subgraph 窗口1 [00:00 - 01:00 限制100]
R1[请求1-100] -->|通过| O1[处理]
R2[请求101] -->|超限| X1[拒绝]
end
subgraph 窗口2 [01:00 - 02:00 限制100]
R3[请求1-100] -->|通过| O2[处理]
end缺点:存在临界点问题——在窗口切换的边界,可能会有两倍的流量通过。例如窗口限制每分钟 100 个请求,但在 00:59 来了 100 个请求,01:01 又来了 100 个请求,这两秒内实际通过了 200 个请求。
package ratelimit
import (
"sync"
"time"
)
// ============================================================
// 固定窗口(Fixed Window)限流算法实现
// ============================================================
// FixedWindow 固定窗口限流器
type FixedWindow struct {
mu sync.Mutex
limit int // 每个窗口允许的最大请求数
window time.Duration // 窗口大小(如 1 分钟)
count int // 当前窗口内的请求计数
windowStart time.Time // 当前窗口的起始时间
}
// NewFixedWindow 创建一个固定窗口限流器
// 参数 limit 是每个窗口允许的最大请求数
// 参数 window 是窗口大小(如 time.Minute 表示 1 分钟)
func NewFixedWindow(limit int, window time.Duration) *FixedWindow {
return &FixedWindow{
limit: limit, // 例如 100 表示每分钟最多 100 个请求
window: window, // 例如 time.Minute
windowStart: time.Now(), // 窗口从当前时间开始
}
}
// Allow 尝试通过限流
func (fw *FixedWindow) Allow() bool {
fw.mu.Lock()
defer fw.mu.Unlock()
now := time.Now()
// 第一步:检查是否需要切换到新窗口
if now.Sub(fw.windowStart) >= fw.window {
// 窗口已过期,开启新窗口
fw.windowStart = now
fw.count = 0
}
// 第二步:检查当前窗口的请求计数
if fw.count < fw.limit {
fw.count++ // 计数加一
return true
}
// 超过限制,拒绝请求
return false
}
3.4 滑动窗口(Sliding Window)
原理:从当前时间开始,往前回溯一段时间,只能处理一定数量的请求。滑动窗口的核心是:这个窗口永远以当前时间戳为准,往前回溯。
flowchart LR
subgraph 滑动窗口
direction TB
P["过去 ←─────────────────→ 现在"]
W["|←── 窗口大小(如1分钟)──→|"]
C["当前时刻往前回溯1分钟内只能N个请求"]
end与固定窗口的对比:
| 特性 | 固定窗口 | 滑动窗口 |
|---|---|---|
| 窗口位置 | 固定在整点开始 | 以当前时刻为终点 |
| 临界点问题 | 有(窗口边界可能通过 2 倍流量) | 无(窗口持续移动) |
| 限流效果 | 不够平滑 | 更加平滑 |
| 实现复杂度 | 简单 | 稍复杂 |
package ratelimit
import (
"sync"
"time"
)
// ============================================================
// 滑动窗口(Sliding Window)限流算法实现
// 使用滑动日志方式实现
// ============================================================
// SlidingWindow 滑动窗口限流器
type SlidingWindow struct {
mu sync.Mutex
limit int // 窗口内允许的最大请求数
window time.Duration // 窗口大小(如 1 分钟)
timestamps []time.Time // 记录每个请求的时间戳
}
// NewSlidingWindow 创建一个滑动窗口限流器
// 参数 limit 是窗口内允许的最大请求数
// 参数 window 是窗口大小(如 time.Minute)
func NewSlidingWindow(limit int, window time.Duration) *SlidingWindow {
return &SlidingWindow{
limit: limit,
window: window,
timestamps: make([]time.Time, 0, limit),
}
}
// Allow 尝试通过限流
func (sw *SlidingWindow) Allow() bool {
sw.mu.Lock()
defer sw.mu.Unlock()
now := time.Now()
// 窗口起始时间 = 当前时间 - 窗口大小
windowStart := now.Add(-sw.window)
// 第一步:清理过期的请求记录(在窗口之外的)
// 从前往后找,把所有早于窗口起始时间的记录删除
idx := 0
for idx < len(sw.timestamps) && sw.timestamps[idx].Before(windowStart) {
idx++
}
sw.timestamps = sw.timestamps[idx:]
// 第二步:检查窗口内的请求数量
if len(sw.timestamps) < sw.limit {
// 未超过限制,记录当前请求的时间戳
sw.timestamps = append(sw.timestamps, now)
return true
}
// 超过限制,拒绝请求
return false
}
滑动窗口的核心区别:固定窗口的窗口边界是固定的(如整点),而滑动窗口的窗口边界是随着当前时间持续移动的,因此限流更加平滑。
3.5 两种窗口对比
固定窗口在窗口切换的瞬间可能允许双倍流量通过(临界点问题),而滑动窗口因为窗口持续移动,不会有这个问题。但滑动窗口的实现稍微复杂一些,需要记录每个请求的时间戳。
四、熔断与降级实现
4.1 熔断器状态机
熔断器的核心是一个状态机,包含三个状态:
flowchart LR
C[Closed
正常状态] -->|错误率超阈值| O[Open
熔断状态]
O -->|等待冷却时间| H[Half-Open
半开状态]
H -->|试探请求成功| C
H -->|试探请求失败| O| 状态 | 说明 |
|---|---|
| Closed(关闭) | 正常状态,请求正常通过。同时统计错误率,当错误率超过阈值时切换到 Open |
| Open(打开) | 熔断状态,所有请求直接被拒绝。等待冷却时间后切换到 Half-Open |
| Half-Open(半开) | 放行少量试探请求。如果成功则回到 Closed,失败则回到 Open |
4.2 熔断器 Go 实现
package circuitbreaker
import (
"errors"
"sync"
"time"
)
// ============================================================
// 熔断器(Circuit Breaker)实现
// ============================================================
// State 熔断器状态
type State int
const (
StateClosed State = iota // 关闭状态:正常处理请求
StateOpen // 打开状态:拒绝所有请求
StateHalfOpen // 半开状态:放行少量试探请求
)
// CircuitBreaker 熔断器
type CircuitBreaker struct {
mu sync.Mutex
state State // 当前状态
failureThreshold int // 失败阈值(连续失败多少次触发熔断)
failureCount int // 当前失败计数
successThreshold int // 半开状态下成功多少次恢复
successCount int // 半开状态下成功计数
cooldown time.Duration // 熔断冷却时间
lastFailureTime time.Time // 上次失败时间
}
// NewCircuitBreaker 创建一个熔断器
// 参数 failureThreshold 是连续失败多少次触发熔断
// 参数 successThreshold 是半开状态下成功多少次恢复到 Closed
// 参数 cooldown 是熔断后的冷却时间
func NewCircuitBreaker(failureThreshold, successThreshold int, cooldown time.Duration) *CircuitBreaker {
return &CircuitBreaker{
state: StateClosed, // 初始状态为关闭
failureThreshold: failureThreshold, // 例如 5 次
successThreshold: successThreshold, // 例如 3 次
cooldown: cooldown, // 例如 30 秒
}
}
// Allow 检查是否允许请求通过
// 返回 nil 表示允许,返回 error 表示被熔断
func (cb *CircuitBreaker) Allow() error {
cb.mu.Lock()
defer cb.mu.Unlock()
now := time.Now()
switch cb.state {
case StateClosed:
// 关闭状态:允许所有请求
return nil
case StateOpen:
// 打开状态:检查冷却时间是否已过
if now.Sub(cb.lastFailureTime) >= cb.cooldown {
// 冷却时间已过,切换到半开状态
cb.state = StateHalfOpen
cb.successCount = 0
return nil
}
// 冷却时间未过,拒绝请求
return errors.New("circuit breaker is open")
case StateHalfOpen:
// 半开状态:允许少量请求通过
return nil
}
return nil
}
// OnSuccess 记录请求成功
func (cb *CircuitBreaker) OnSuccess() {
cb.mu.Lock()
defer cb.mu.Unlock()
switch cb.state {
case StateClosed:
// 关闭状态下成功,重置失败计数
cb.failureCount = 0
case StateHalfOpen:
// 半开状态下成功,增加成功计数
cb.successCount++
if cb.successCount >= cb.successThreshold {
// 达到成功阈值,恢复到关闭状态
cb.state = StateClosed
cb.failureCount = 0
cb.successCount = 0
}
}
}
// OnFailure 记录请求失败
func (cb *CircuitBreaker) OnFailure() {
cb.mu.Lock()
defer cb.mu.Unlock()
cb.lastFailureTime = time.Now()
switch cb.state {
case StateClosed:
// 关闭状态下失败,增加失败计数
cb.failureCount++
if cb.failureCount >= cb.failureThreshold {
// 达到失败阈值,切换到打开状态
cb.state = StateOpen
}
case StateHalfOpen:
// 半开状态下失败,重新切换到打开状态
cb.state = StateOpen
cb.failureCount = 0
cb.successCount = 0
}
}
4.3 gRPC 限流拦截器实现
package interceptor
import (
"context"
"sync"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)
// ============================================================
// gRPC 服务端限流拦截器
// 使用令牌桶算法对每个方法进行限流
// ============================================================
// RateLimitInterceptor 限流拦截器
// 对每个 gRPC 方法维护一个独立的令牌桶
type RateLimitInterceptor struct {
mu sync.Mutex
limiters map[string]*TokenBucket // 方法名 -> 令牌桶
rate float64 // 令牌生成速率
capacity float64 // 桶容量
}
// NewRateLimitInterceptor 创建限流拦截器
// 参数 rate 是每秒允许的请求数
// 参数 capacity 是令牌桶容量(允许的突发请求数)
func NewRateLimitInterceptor(rate, capacity float64) *RateLimitInterceptor {
return &RateLimitInterceptor{
limiters: make(map[string]*TokenBucket),
rate: rate,
capacity: capacity,
}
}
// getLimiter 获取或创建某个方法的令牌桶
// 不同的方法可以使用不同的令牌桶,互不影响
func (r *RateLimitInterceptor) getLimiter(method string) *TokenBucket {
r.mu.Lock()
defer r.mu.Unlock()
if limiter, ok := r.limiters[method]; ok {
return limiter
}
// 为新方法创建令牌桶
limiter := NewTokenBucket(r.rate, r.capacity)
r.limiters[method] = limiter
return limiter
}
// ServerInterceptor 返回 gRPC 服务端拦截器
// 在每个请求处理前检查是否超过限流阈值
func (r *RateLimitInterceptor) ServerInterceptor() grpc.UnaryServerInterceptor {
return func(
ctx context.Context, // 请求上下文
req interface{}, // 请求参数
info *grpc.UnaryServerInfo, // 服务端调用信息,包含方法名
handler grpc.UnaryHandler, // 实际的业务处理函数
) (interface{}, error) {
// 第一步:获取该方法的令牌桶
limiter := r.getLimiter(info.FullMethod)
// 第二步:检查是否允许通过
if !limiter.Allow() {
// 超过限流,返回 ResourceExhausted 错误
// gRPC 标准状态码中,ResourceExhausted 表示资源不足
return nil, status.Errorf(codes.ResourceExhausted,
"rate limit exceeded for method %s", info.FullMethod)
}
// 第三步:调用实际业务逻辑
return handler(ctx, req)
}
}
4.4 gRPC 熔断拦截器实现
package interceptor
import (
"context"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)
// ============================================================
// gRPC 客户端熔断拦截器
// ============================================================
// CircuitBreakerInterceptor 熔断拦截器
type CircuitBreakerInterceptor struct {
breaker *CircuitBreaker
}
// NewCircuitBreakerInterceptor 创建熔断拦截器
func NewCircuitBreakerInterceptor(breaker *CircuitBreaker) *CircuitBreakerInterceptor {
return &CircuitBreakerInterceptor{breaker: breaker}
}
// ClientInterceptor 返回 gRPC 客户端拦截器
// 在每次调用前检查熔断状态,调用后记录成功/失败
func (c *CircuitBreakerInterceptor) ClientInterceptor() grpc.UnaryClientInterceptor {
return func(
ctx context.Context, // 请求上下文
method string, // 方法名
req, reply interface{}, // 请求和响应
cc *grpc.ClientConn, // 客户端连接
invoker grpc.UnaryInvoker, // 实际调用函数
opts ...grpc.CallOption, // 调用选项
) error {
// 第一步:检查熔断器是否允许请求通过
if err := c.breaker.Allow(); err != nil {
// 熔断器打开,直接返回错误
return status.Error(codes.Unavailable, "circuit breaker is open")
}
// 第二步:发起实际调用
err := invoker(ctx, method, req, reply, cc, opts...)
// 第三步:根据调用结果更新熔断器状态
if err != nil {
c.breaker.OnFailure() // 调用失败,记录失败
} else {
c.breaker.OnSuccess() // 调用成功,记录成功
}
return err
}
}
4.5 降级实现
降级的核心是为业务准备快路径和慢路径:
package degradation
import (
"context"
"google.golang.org/grpc"
)
// ============================================================
// 降级拦截器实现
// 当服务不可用时,走简易路径返回默认值
// ============================================================
// degradeKey 用于在 context 中标记降级请求
type degradeKey struct{}
// WithDegrade 在 context 中设置降级标记
// 被标记的请求会走快路径(简易逻辑)
func WithDegrade(ctx context.Context) context.Context {
return context.WithValue(ctx, degradeKey{}, true)
}
// IsDegrade 检查是否为降级请求
func IsDegrade(ctx context.Context) bool {
v, ok := ctx.Value(degradeKey{}).(bool)
return ok && v
}
// DegradeInterceptor 降级拦截器
// 参数 fallback 是降级时执行的简易逻辑
// 当原始调用失败时,执行 fallback 返回默认值
func DegradeInterceptor(fallback func(ctx context.Context, method string, req interface{}) (interface{}, error)) grpc.UnaryClientInterceptor {
return func(
ctx context.Context,
method string,
req, reply interface{},
cc *grpc.ClientConn,
invoker grpc.UnaryInvoker,
opts ...grpc.CallOption,
) error {
// 第一步:检查是否已经被标记为降级请求
if IsDegrade(ctx) {
// 已经是降级请求,直接走快路径
result, err := fallback(ctx, method, req)
if err != nil {
return err
}
// 将 fallback 的结果复制到 reply
// 实际实现中需要使用反射或 protobuf 的 Merge 方法
return nil
}
// 第二步:正常调用
err := invoker(ctx, method, req, reply, cc, opts...)
if err != nil {
// 第三步:调用失败,走降级逻辑
result, ferr := fallback(ctx, method, req)
if ferr == nil && result != nil {
// 降级成功,返回默认值
_ = result // 实际实现中需要将 result 复制到 reply
return nil
}
// 降级也失败了,返回原始错误
return err
}
return nil
}
}
降级在面试中的回答:业务分成了快路径和慢路径两种。慢路径很消耗资源,是正常业务逻辑;快路径可以是直接返回默认值,也可以是存储数据后面异步处理。不降级时先走快路径再走慢路径,降级时只走快路径。
五、客户端限流与服务端限流
5.1 客户端限流 vs 服务端限流
| 维度 | 服务端限流 | 客户端限流 |
|---|---|---|
| 位置 | 在服务端拦截器中执行 | 在客户端拦截器中执行 |
| 优点 | 精确控制服务端负载 | 避免无意义的网络请求 |
| 缺点 | 请求已经到达服务端 | 不同客户端各自限流,合在一起可能超量 |
| 使用频率 | 常用 | 较少使用 |
客户端限流用得少,主要是因为很可能不同客户端上单独限流了,结果合在一起却超过了服务器处理能力。
5.2 单机限流与集群限流
前面讨论的都是单机限流,在微服务框架下还可以考虑对集群进行限流。
集群限流特征:
- 非常接近网关限流
- 集群限流主要依赖于在不同的实例之间同步阈值、当前请求数,目前使用 Redis 的比较多(高并发场景)
5.3 基于 Redis 的集群限流
固定窗口的 Redis 实现
package ratelimit
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
// ============================================================
// 基于 Redis 的集群限流实现(固定窗口)
// ============================================================
// RedisFixedWindow 基于 Redis 的固定窗口集群限流器
type RedisFixedWindow struct {
client *redis.Client // Redis 客户端
limit int // 每个窗口允许的最大请求数
window time.Duration // 窗口大小
}
// NewRedisFixedWindow 创建 Redis 集群限流器
func NewRedisFixedWindow(client *redis.Client, limit int, window time.Duration) *RedisFixedWindow {
return &RedisFixedWindow{
client: client,
limit: limit, // 例如 1000 表示每分钟整个集群最多 1000 个请求
window: window, // 例如 time.Minute
}
}
// Allow 尝试通过限流
// 使用 Redis 的 INCR 命令原子递增计数器
func (r *RedisFixedWindow) Allow(ctx context.Context, key string) (bool, error) {
// 第一步:计算当前窗口的 Redis key
// 使用时间戳取整来确保同一窗口内的 key 相同
// 例如窗口为 1 分钟,则 00:01:30 和 00:01:59 的 key 相同
now := time.Now()
windowStart := now.Truncate(r.window) // 截断到窗口边界
redisKey := fmt.Sprintf("ratelimit:%s:%d", key, windowStart.Unix())
// 第二步:使用 Lua 脚本原子操作
// INCR 和 EXPIRE 必须在同一个原子操作中执行
// 否则可能出现 INCR 成功但 EXPIRE 失败导致 key 永不过期
script := `
local count = redis.call('INCR', KEYS[1])
if count == 1 then
redis.call('EXPIRE', KEYS[1], ARGV[1])
end
return count
`
// 窗口过期时间(秒)
ttl := int64(r.window.Seconds())
// 执行 Lua 脚本
result, err := r.client.Eval(ctx, script, []string{redisKey}, ttl).Int()
if err != nil {
return false, err
}
// 第三步:判断是否超过限制
return result <= r.limit, nil
}
滑动窗口的 Redis 实现
package ratelimit
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
// ============================================================
// 基于 Redis 的集群限流实现(滑动窗口)
// 使用 Sorted Set(有序集合)实现
// ============================================================
// RedisSlidingWindow 基于 Redis 的滑动窗口集群限流器
type RedisSlidingWindow struct {
client *redis.Client // Redis 客户端
limit int // 窗口内允许的最大请求数
window time.Duration // 窗口大小
}
// NewRedisSlidingWindow 创建 Redis 滑动窗口集群限流器
func NewRedisSlidingWindow(client *redis.Client, limit int, window time.Duration) *RedisSlidingWindow {
return &RedisSlidingWindow{
client: client,
limit: limit,
window: window,
}
}
// Allow 尝试通过限流
// 使用 Redis 的 Sorted Set 记录请求时间戳
func (r *RedisSlidingWindow) Allow(ctx context.Context, key string) (bool, error) {
now := time.Now()
windowStart := now.Add(-r.window) // 窗口起始时间
// Redis key
redisKey := fmt.Sprintf("sliding_window:%s", key)
// 使用 Lua 脚本保证原子性
// 步骤:
// 1. 移除窗口外的旧记录
// 2. 检查当前窗口内的请求数
// 3. 如果未超限,添加当前请求的时间戳
script := `
-- 第一步:移除窗口外的旧记录
redis.call('ZREMRANGEBYSCORE', KEYS[1], 0, ARGV[1])
-- 第二步:获取当前窗口内的请求数
local count = redis.call('ZCARD', KEYS[1])
-- 第三步:判断是否超过限制
if count < tonumber(ARGV[3]) then
-- 未超限,添加当前请求
redis.call('ZADD', KEYS[1], ARGV[2], ARGV[2])
-- 设置 key 的过期时间,防止内存浪费
redis.call('EXPIRE', KEYS[1], ARGV[4])
return 1
else
return 0
end
`
// 参数说明:
// ARGV[1] = 窗口起始时间(用于移除旧记录)
// ARGV[2] = 当前时间戳(作为 score 和 member)
// ARGV[3] = 限制数量
// ARGV[4] = 过期时间(秒)
result, err := r.client.Eval(ctx, script,
[]string{redisKey},
windowStart.UnixNano(),
now.UnixNano(),
r.limit,
int64(r.window.Seconds()),
).Int()
if err != nil {
return false, err
}
return result == 1, nil
}
本质上,Redis 集群限流就是把单机限流的算法用 Lua 脚本再写了一遍,利用 Redis 的原子性来保证集群维度的限流正确性。
5.4 限流对象
除了单机限流和集群限流,还有其他维度的限流:
| 维度 | 说明 | 示例 |
|---|---|---|
| 接口维度 | 对每个接口单独限流 | 用户接口 100 QPS,订单接口 50 QPS |
| 方法维度 | 对每个方法单独限流 | GetUser 方法 100 QPS |
| 用户维度 | 针对特定用户限流 | 用户 A 最多 10 QPS |
| IP 维度 | 针对来源 IP 限流 | 同一 IP 不能频繁登录 |
业务相关限流一般针对安全、业务价值。大体上的逻辑就是牺牲不重要的来保护重要的。也可以考虑跨服务限流——如果服务很重要,可以将不重要的服务流量摘掉,腾出资源来保护核心服务。
5.5 拒绝策略
在实际中,限流不仅仅是返回错误,还有多种拒绝策略:
| 策略 | 说明 |
|---|---|
| 标记简易路径 | 标记请求为限流请求,后续业务走简易路径(如缓存未命中直接返回,不查数据库) |
| 返回固定响应 | 直接在拦截器中返回固定的默认响应 |
| 缓存请求 | 将请求存到数据库或 Redis 中,后面再取出重做 |
| 转为异步 | 拦截器返回类似 202 响应(请求已接收),后续由调度器异步执行 |
| 转发到别的服务器 | 限流的请求转发到其他机器(类似 302 重定向) |
六、可观测性三支柱
可观测性(Observability)是微服务运维的基础,包含三个支柱:
flowchart TB
O[可观测性
Observability] --> M[Metrics
指标监控]
O --> T[Tracing
链路追踪]
O --> L[Logging
日志记录]| 支柱 | 说明 | 关注的问题 |
|---|---|---|
| Metrics | 记录系统运行时的量化指标 | 系统当前健康吗?响应时间多少?错误率多少? |
| Tracing | 记录一个请求在多个服务间的调用链路 | 请求在哪个环节慢了?哪个服务出了错? |
| Logging | 记录离散的事件日志 | 具体发生了什么?请求参数是什么? |
6.1 Metrics(指标监控)
Metrics 记录的是系统运行的量化数据,通常包括:
- 请求总数:QPS(每秒请求数)
- 响应时间:平均响应时间、P99 响应时间
- 错误数:错误总数、错误率
- 请求/响应大小:数据传输量
package metrics
import (
"context"
"time"
"google.golang.org/grpc"
)
// ============================================================
// gRPC Metrics 拦截器
// 记录请求的响应时间和调用结果
// ============================================================
// MetricsInterceptor Metrics 拦截器
type MetricsInterceptor struct {
// 实际项目中,这里会接入 Prometheus 等监控系统
// 我们这里简化演示,用 map 存储统计数据
}
// NewMetricsInterceptor 创建 Metrics 拦截器
func NewMetricsInterceptor() *MetricsInterceptor {
return &MetricsInterceptor{}
}
// ServerInterceptor 返回服务端 Metrics 拦截器
// 在请求处理前后记录指标数据
func (m *MetricsInterceptor) ServerInterceptor() grpc.UnaryServerInterceptor {
return func(
ctx context.Context, // 请求上下文
req interface{}, // 请求参数
info *grpc.UnaryServerInfo, // 调用信息
handler grpc.UnaryHandler, // 业务处理函数
) (interface{}, error) {
// 第一步:记录请求开始时间
start := time.Now()
// 第二步:调用实际业务逻辑
resp, err := handler(ctx, req)
// 第三步:计算响应时间
duration := time.Since(start)
// 第四步:记录指标
// 实际项目中,这些数据会发送到 Prometheus 等监控系统
// 这里简化演示
m.recordMetrics(info.FullMethod, duration, err)
return resp, err
}
}
// recordMetrics 记录指标数据
// method: 方法名,如 "/user.UserService/GetUser"
// duration: 响应时间
// err: 调用错误(nil 表示成功)
func (m *MetricsInterceptor) recordMetrics(method string, duration time.Duration, err error) {
// 在实际项目中,这些数据会被 Prometheus 采集
// 例如:
// - grpc_request_total{method="GetUser", status="success"} ++
// - grpc_request_duration_seconds{method="GetUser"} observe(duration)
// 这里简化输出
status := "success"
if err != nil {
status = "error"
}
// 实际实现中,不要在拦截器里打印日志(影响性能)
// 这里仅作演示
_ = method
_ = duration
_ = status
}
6.2 Tracing(链路追踪)
在微服务架构中,一个请求可能经过多个服务。Tracing 的作用是将这些调用串联起来,形成完整的调用链路。
flowchart LR
C[客户端] -->|TraceID: abc123| A[服务A]
A -->|传递 TraceID| B[服务B]
B -->|传递 TraceID| D[服务C]
D -->|返回| B
B -->|返回| A
A -->|返回| C核心概念:
| 概念 | 说明 |
|---|---|
| TraceID | 全局唯一标识一次完整的调用链路 |
| SpanID | 标识链路中的一个节点(一次 RPC 调用) |
| ParentSpanID | 父节点的 SpanID,用于构建调用树 |
| Extract | 从请求中提取链路元数据(如从 gRPC metadata 中提取 TraceID) |
| Inject | 将链路元数据注入到请求中(如将 TraceID 写入 gRPC metadata) |
package tracing
import (
"context"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/metadata"
)
// ============================================================
// gRPC Tracing 拦截器
// 将请求在多个服务间的调用串联起来
// ============================================================
// 链路追踪的 metadata key
const (
TraceIDKey = "x-trace-id"
SpanIDKey = "x-span-id"
ParentSpanIDKey = "x-parent-span-id"
)
// generateID 生成唯一 ID(简化版,实际使用 UUID 或类似方案)
func generateID() string {
return time.Now().Format("20060102150405.000000")
}
// ServerTracingInterceptor 服务端链路追踪拦截器
// 从 gRPC metadata 中提取 TraceID,创建新的 SpanID
func ServerTracingInterceptor() grpc.UnaryServerInterceptor {
return func(
ctx context.Context,
req interface{},
info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler,
) (interface{}, error) {
// 第一步:从 metadata 中提取链路信息(Extract)
// 客户端在发起请求时会将 TraceID 放入 metadata
md, ok := metadata.FromIncomingContext(ctx)
var traceID string
var parentSpanID string
if ok {
// 提取 TraceID
if values := md.Get(TraceIDKey); len(values) > 0 {
traceID = values[0]
}
// 提取父 SpanID
if values := md.Get(SpanIDKey); len(values) > 0 {
parentSpanID = values[0]
}
}
// 如果没有 TraceID,说明是链路起点,生成新的
if traceID == "" {
traceID = generateID()
}
// 第二步:生成当前节点的 SpanID
spanID := generateID()
// 第三步:记录 Span 信息
// 实际项目中,这些信息会发送到 Jaeger、Zipkin 等 tracing 系统
start := time.Now()
// 将链路信息存入 context,供后续使用
ctx = context.WithValue(ctx, "traceID", traceID)
ctx = context.WithValue(ctx, "spanID", spanID)
// 第四步:调用业务逻辑
resp, err := handler(ctx, req)
// 第五步:记录 Span 的结束信息
duration := time.Since(start)
// 实际项目中,会将 span 信息上报到 tracing 系统:
// span.SetTag("method", info.FullMethod)
// span.SetTag("duration", duration)
// span.SetTag("error", err != nil)
// span.Finish()
_ = duration
_ = err
return resp, err
}
}
// ClientTracingInterceptor 客户端链路追踪拦截器
// 将 TraceID 注入到 gRPC metadata 中(Inject)
func ClientTracingInterceptor() grpc.UnaryClientInterceptor {
return func(
ctx context.Context,
method string,
req, reply interface{},
cc *grpc.ClientConn,
invoker grpc.UnaryInvoker,
opts ...grpc.CallOption,
) error {
// 第一步:从 context 中获取链路信息
traceID, _ := ctx.Value("traceID").(string)
if traceID == "" {
// 如果没有 TraceID,生成新的
traceID = generateID()
}
// 生成当前调用的 SpanID
spanID := generateID()
// 第二步:将链路信息注入到 metadata 中(Inject)
md := metadata.Pairs(
TraceIDKey, traceID, // 传递 TraceID
SpanIDKey, spanID, // 当前的 SpanID 成为下游的 ParentSpanID
)
// 将 metadata 附加到 context
ctx = metadata.AppendToOutgoingContext(ctx,
TraceIDKey, traceID,
SpanIDKey, spanID,
)
// 第三步:发起调用
return invoker(ctx, method, req, reply, cc, opts...)
}
}
Kratos 的 tracing 比较有趣的是记录了请求和响应的大小,而不仅仅是响应时间。
6.3 Logging(日志记录)
Logging 记录的是离散的事件日志,通常包括请求参数、响应结果、错误信息等。
package logging
import (
"context"
"log"
"time"
"google.golang.org/grpc"
)
// ============================================================
// gRPC Logging 拦截器
// 记录请求和响应的详细信息
// ============================================================
// LoggingInterceptor 日志拦截器
type LoggingInterceptor struct {
logger *log.Logger // 日志记录器
}
// NewLoggingInterceptor 创建日志拦截器
func NewLoggingInterceptor(logger *log.Logger) *LoggingInterceptor {
return &LoggingInterceptor{logger: logger}
}
// ServerInterceptor 返回服务端日志拦截器
func (l *LoggingInterceptor) ServerInterceptor() grpc.UnaryServerInterceptor {
return func(
ctx context.Context,
req interface{},
info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler,
) (interface{}, error) {
start := time.Now()
// 记录请求信息
// 注意:记录请求参数要考虑两个问题:
// 1. 请求可能很大(占用日志存储空间)
// 2. 请求可能包含敏感数据(如住址、手机号码等)
l.logger.Printf("[REQ] method=%s req=%v", info.FullMethod, req)
// 调用业务逻辑
resp, err := handler(ctx, req)
// 记录响应信息
duration := time.Since(start)
if err != nil {
l.logger.Printf("[ERR] method=%s duration=%s err=%v",
info.FullMethod, duration, err)
} else {
l.logger.Printf("[RSP] method=%s duration=%s",
info.FullMethod, duration)
}
return resp, err
}
}
记录请求参数要注意两个问题:
- 请求很大——占用大量日志存储空间
- 请求包含敏感数据——如住址、手机号码等,需要脱敏处理
6.4 各框架的可观测性对比
| 框架 | Metrics | Tracing | Logging |
|---|---|---|---|
| Kratos | 支持分错误码观测,记录响应时间 | 利用错误传递机制记录错误码,记录请求/响应大小 | 记录请求参数 |
| Dubbo-go | 记录响应时间等 | 将链路串联在一起,额外记录错误信息 | 记录调用信息 |
| go-micro | 支持 Prometheus,观测响应时间、请求总数、错误总数 | 支持多种 tracing 工具(如 OpenTelemetry) | 记录调用日志 |
七、基于可观测性的服务治理
7.1 静态策略与动态策略
服务治理策略分为静态和动态两类:
| 类型 | 负载均衡算法 | 限流算法 |
|---|---|---|
| 静态策略 | 随机、哈希、轮询 | 固定窗口、滑动窗口、漏桶、令牌桶 |
| 动态策略 | 最少连接数、最少活跃请求数 | BBR |
7.2 静态策略的缺陷
静态算法依赖于程序员的个人经验。例如在限流中,有三种方式确定限流的阈值:
- 基于压测结果——通过压力测试找到系统的承受能力
- 基于可观测性数据——根据线上监控数据调整
- 基于源码分析——分析代码逻辑推算
对于负载均衡,有一些基本假设:
- 所有的请求都消耗一样的资源
- 大量的请求
如果按照
user_id % 3进行哈希负载均衡,并且如果user_id 1是热门用户(例如大 UP 主),那么节点 1 就可能过载。这是因为静态算法不考虑请求的实际资源消耗。
7.3 动态策略的缺陷
动态策略试图寻找一些指标来判断节点的健康程度:
- 最少连接数:使用连接数,连接数越少认为越健康
- 最少活跃请求数:使用当前正在处理的请求数量,越少认为越健康
- 最快响应时间:使用响应时间,越快认为越健康
核心问题:微服务框架本身并不具备全局信息。例如,节点 1 有 120 个连接,节点 2 只有 30 个连接,但客户端只知道自己到这些节点的连接情况,不知道其他客户端的连接情况。所以客户端最终可能会做出错误的选择。
动态调整权重类的负载均衡算法,就是试图通过增加权重或降低权重来表达一个节点的健康程度。例如在超时的时候降低权重,而在拿到了响应之后增加权重。这一类的算法在计算权重的时候需要非常小心。
7.4 可利用的指标
大多数情况下,静态策略就运作良好。而动态策略则可以考虑使用:
| 指标类型 | 具体指标 |
|---|---|
| 硬件/环境指标 | CPU 利用率、内存利用率、GC 时间 |
| 服务指标 | 响应时间、错误率、超时率 |
这些指标并不是独立的,相互之间有影响。例如 CPU 利用率高可能导致响应时间变长,响应时间变长可能导致超时率上升。
7.5 基于可观测性的服务治理基本思路
基本思路(以负载均衡为例):
flowchart LR
A[可观测性平台
从所有节点采集指标] --> B[客户端拉取指标
或平台推送指标]
B --> C[客户端使用指标
执行负载均衡]- 可观测性平台从所有的节点采集指标
- 客户端拉取指标(或者可观测性平台推送指标)
- 客户端使用这些指标来执行负载均衡
开发者可以根据指标和业务特征来设计负载均衡算法。
7.6 指标的时间敏感性
大多数指标都是时间敏感的——你拿到的指标可能是几秒甚至几分钟前的,不能准确反映当前状态。
为了规避采集指标的延时问题,有两种方案:
方案一:响应附带指标
在返回响应的时候将指标一起返回。在高并发的环境下,这种策略可以解决健康检查引起的网络性能问题。
- 优点:指标是实时的
- 缺点:
- 一些 RPC 协议不支持从服务端返回这类数据(协议中没有预留字段)
- 如果客户端长期没有发送请求,它持有的数据都是很久以前的
方案二:可观测性平台采集
通过可观测性平台统一采集和分发指标。
- 优点:覆盖全面,不依赖请求
- 缺点:有采集延迟
综合方案
两者结合能够有效利用两者的优点而且规避缺点。在这种机制之下,客户端可以从节点 1 中知晓它上面有 120 个连接,而节点 2 上面只有 30 个连接,从而做出正确的选择。
7.7 服务端治理与客户端治理
相比之下,服务端的治理要简单很多,不需要考虑那么复杂的时间敏感性问题。因为它本身就有自己的全部信息,而且是实时信息。例如可以根据自身统计的响应时间、错误率、CPU 等来实时计算是否需要限流。
7.8 第三方组件指标
如果寻求一个万无一失的方案,还需要采集第三方组件的指标,例如系统使用的 Redis 集群、数据库集群等。服务端的治理(如限流)也可以利用第三方采集的指标,尤其是大部分系统的核心组件——数据库的指标,非常具有参考意义。
7.9 利用指标
如何使用这些指标?答案是:水无常势,兵无定型,取决于你的业务特征。
大多数情况下,选择几个关键的指标就可以了:
- CPU + 内存 + 网络 IO
- 响应时间或错误率等
八、总结与面试要点
8.1 限流、熔断、降级总结
| 概念 | 本质 | 触发条件 | 处理方式 |
|---|---|---|---|
| 限流 | 只允许一部分请求被处理 | 请求量超过阈值 | 超限请求被拒绝或排队 |
| 熔断 | 拒绝所有请求 | 错误率超过阈值 | 直接返回错误 |
| 降级 | 走简易逻辑 | 服务不可用或负载高 | 返回默认值或走快路径 |
三者之间其实并不是泾渭分明的,本质上只是处理策略上稍微有点区别。从接口设计和算法层面来说,它们也非常接近。
8.2 限流算法对比
| 算法 | 类型 | 优点 | 缺点 |
|---|---|---|---|
| 令牌桶 | 静态 | 允许突发流量 | 需要选择合适的速率和容量 |
| 漏桶 | 静态 | 输出速率严格固定 | 不允许突发流量 |
| 固定窗口 | 静态 | 实现简单 | 有临界点问题 |
| 滑动窗口 | 静态 | 限流平滑 | 实现稍复杂 |
8.3 面试要点
可用性方面
- 限流的几个算法? 掌握令牌桶、漏桶、固定窗口、滑动窗口的原理和代码实现
- 滑动窗口和固定窗口的区别? 滑动窗口的限流更加平滑,因为窗口是在持续移动的
- 令牌桶和漏桶的区别? 两者效果一样,令牌桶允许突发流量,漏桶输出速率严格固定
- 什么是降级? 在业务层面上准备快慢两条路径。不降级时执行慢路径(正常逻辑),降级时走快路径(默认值或异步处理)
- 什么是熔断? 在系统故障时拒绝新的请求。如果不是拒绝新请求而是走简易逻辑,也可以说是降级
- 什么是限流? 在一定时间段内只允许一部分请求被处理,其余请求被拒绝或降级
- 触发降级和熔断后怎么恢复? 常见做法是过一段时间后直接退出降级/熔断状态;高级做法是退出前先试探一下,放过去少部分请求
- 三者之间的联系和区别? 三者并不是泾渭分明的,本质上只是处理策略上的区别
- 可以针对什么来限流? 单机限流、集群限流、业务限流(用户限流、IP 限流)
可观测性方面
- 可观测性的基本概念——Metrics、Tracing、Logging 三大支柱
- 微服务应该采集哪些指标? 注意两个方面:绝对值和趋势。例如平均响应时间,以及平均响应时间的变化趋势
- 怎么把可观测性和服务治理结合起来? 大部分面试官想不到可以将可观测性平台和服务治理结合起来。核心思路是:从可观测性平台采集指标,用指标驱动负载均衡、限流等治理决策
- 告警系统怎么做? 利用观测到的数据,设定各种告警阈值(如错误率超过 1% 就告警),考虑告警手段(邮件、即时通讯、电话)。监控和告警一般是一体的
总结
本教程从 AOP 方案对比出发,详细讲解了微服务框架中可用性(熔断、限流、降级)和可观测性(Metrics、Tracing、Logging)的核心概念和代码实现,最后探讨了如何将可观测性数据与服务治理结合起来。
关键知识点回顾:
- AOP 是可用性和可观测性的基础——通过拦截器在不侵入业务代码的前提下增加治理逻辑
- 熔断、限流、降级本质相同——都是故障检测 + 故障处理 + 故障恢复,区别在于处理策略
- 限流算法中令牌桶和漏桶效果一样,滑动窗口比固定窗口更平滑
- 熔断器是一个三状态状态机:Closed -> Open -> Half-Open
- 集群限流通过 Redis + Lua 脚本实现,保证集群维度的原子性
- 可观测性三支柱:Metrics(指标)、Tracing(链路追踪)、Logging(日志)
- 链路追踪通过 TraceID 将多个服务的调用串联起来,核心是 Extract 和 Inject
- 基于可观测性的服务治理是将监控数据与负载均衡、限流等治理策略结合
- 指标的时间敏感性是动态治理的核心挑战,综合方案(响应附带 + 平台采集)是最优解
- 服务端治理比客户端治理简单,因为服务端拥有自己的全部实时信息
自测题与动手练习
自测题(合上书能答出来,才算懂):
- 熔断、限流、降级三者本质是否不同?它们之间最核心的区别是什么?
- 滑动窗口相比固定窗口解决了什么问题?为什么固定窗口在窗口边界可能放行近 2 倍流量?
- 令牌桶和漏桶"效果一样"指的是什么?令牌桶相对漏桶多出的能力是什么?
- 熔断器三个状态(Closed / Open / Half-Open)之间如何转换?引入 Half-Open 半开状态是为了解决什么?
- 可观测性三支柱分别回答什么问题?TraceID 和 SpanID(含 ParentSpanID)各自的职责是什么?
动手练习(建议真做一遍):
- 把本文的令牌桶
Allow()改造成"阻塞等待"版本(拿不到令牌就time.Sleep一小会儿重试,直到拿到或整体超时),对比"直接拒绝"与"阻塞排队"两种策略在突发流量下的表现。 - 把 6.1 的
MetricsInterceptor真正接入 Prometheus(用官方client_golang计数器与直方图),用grpc_request_duration_seconds看到 QPS 与 P99 曲线。 - 在 5.3 Redis 集群限流的
RedisFixedWindow上,故意把INCR和EXPIRE拆成两条非原子命令,写一个并发脚本复现"key 永不过期导致限流失效"的 bug,再换回 Lua 脚本验证修复。
本章小结
- AOP 是地基:可用性(限流 / 熔断 / 降级)与可观测性(Metrics / Tracing / Logging)都通过 gRPC Interceptor 在不侵入业务代码的前提下织入。
- 三者本质相同:熔断、限流、降级都是"故障检测 → 故障处理 → 故障恢复",差别只在处理策略(拒绝 / 排队 / 走快路径)。
- 限流算法选型:令牌桶与漏桶效果等价,令牌桶允许突发;固定窗口实现简单但有临界点问题;滑动窗口更平滑但需记录时间戳。
- 熔断器是三态状态机:Closed → Open → Half-Open,Half-Open 用少量试探请求决定恢复还是继续熔断。
- 集群限流 = 单机算法 + Redis 原子性:本质是把单机限流用 Lua 脚本在 Redis 上重写一遍。
- 可观测性三支柱:Metrics 看"系统健不健康",Tracing 看"请求卡在哪",Logging 看"具体发生了什么";链路追踪靠 TraceID 串联、Extract / Inject 透传。
- 治理靠数据驱动:可观测性平台采集指标 → 客户端拉取或平台推送 → 驱动负载均衡、限流等决策;服务端治理比客户端治理简单,因为它掌握自身全部实时信息。
下一章我们将进入更贴近生产的环节:把可观测性真正接到 OpenTelemetry 与 ELK 上,构建可落地的日志、追踪与监控体系。