五、可用性和可观测性

2021-02-19T14:21:02+08:00 | 30分钟阅读 | 更新于 2021-02-19T14:21:02+08:00

@

学习目标

类比(先建立直觉):微服务的"可用性"和"可观测性"就像你开了一家连锁餐厅。可用性是"出事时还能接住客人"——后厨着火就挂"暂停营业"(熔断)、客人太多就取号排队(限流)、人手不够就只卖预制菜(降级);可观测性是"装监控、对讲机、点单日志"——出了差错能立刻定位到是哪桌、哪道菜、哪个厨师。本章就是教你怎么把这套"餐厅应急 + 监控"机制用代码落到微服务框架里。

学完本章你应该能够:

  1. 用自己的话讲清熔断、限流、降级三者"本质相同、处理策略不同"的关系,并画出故障处理流程与熔断器状态机。
  2. 讲清令牌桶、漏桶、固定窗口、滑动窗口四种限流算法的原理、差异(尤其是临界点问题、突发流量)与适用场景。
  3. 基于 gRPC Interceptor 写一个可用的限流 / 熔断拦截器,并说清客户端限流 vs 服务端限流、单机限流 vs 集群限流的区别。
  4. 讲清可观测性三支柱(Metrics / Tracing / Logging)的分工,以及 TraceID / SpanID 如何把一次跨服务调用串联成完整链路。
  5. 说出"基于可观测性的服务治理“的基本思路(指标采集 → 客户端拉取 / 平台推送 → 驱动负载均衡与限流),并在面试中讲成一个完整故事。

前置知识

  • 第 1–4 章:网络编程、RPC 协议、服务注册发现、负载均衡与集群容错
  • Go 基础:goroutinesync.Mutexcontexttime
  • gRPC 基本用法与 Interceptor 机制
  • Redis 基础命令:INCR / ZADD / Lua 脚本

本章你会动手做的事

  1. 逐一跑一遍令牌桶 / 漏桶 / 固定窗口 / 滑动窗口的 Allow(),观察高并发下是否如预期拒绝请求。
  2. 写一个 gRPC 服务端限流拦截器(令牌桶),对同一方法做 QPS 限制。
  3. 本地起 Redis,用 Lua 脚本实现一次集群固定窗口限流,压测验证其原子性。

前言:微服务的稳定运行与问题排查

在前面的章节中,我们已经构建了一个完整的微服务框架——从网络编程到 RPC 协议,从服务注册发现到负载均衡和集群容错。但是,一个微服务框架光能跑起来还不够,还需要解决两个关键问题:

  1. 可用性——当服务出现故障或流量激增时,如何保证系统仍然可用?
  2. 可观测性——当系统出问题时,如何快速定位和排查?

本教程涵盖以下主题:

  1. AOP 方案对比——Kratos、Dubbo-go、go-micro 的拦截器设计
  2. 可用性核心概念——熔断、限流、降级的关系与区别
  3. 限流算法详解——令牌桶、漏桶、固定窗口、滑动窗口及 Go 实现
  4. 熔断与降级——gRPC 拦截器实现限流与熔断
  5. 集群限流——基于 Redis 的分布式限流
  6. 可观测性三支柱——Metrics(指标)、Tracing(链路追踪)、Logging(日志)
  7. 基于可观测性的服务治理——如何利用观测数据驱动治理决策
  8. 面试要点总结

一、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

实际上 WrapperCallWrapper 只需要保留一个就可以。通过注册过程来控制作用范围,例如我们在某个调用里面注册 Wrapper,那么只对该调用起效果。

服务端的 Wrapper 定义不太一样,核心就是 HandlerWrapperStreamWrapper,其作用接近于客户端的 WrapperStreamWrapper

设计偏好:个人比较喜欢 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 各框架的可观测性对比

框架MetricsTracingLogging
Kratos支持分错误码观测,记录响应时间利用错误传递机制记录错误码,记录请求/响应大小记录请求参数
Dubbo-go记录响应时间等将链路串联在一起,额外记录错误信息记录调用信息
go-micro支持 Prometheus,观测响应时间、请求总数、错误总数支持多种 tracing 工具(如 OpenTelemetry)记录调用日志

七、基于可观测性的服务治理

7.1 静态策略与动态策略

服务治理策略分为静态和动态两类:

类型负载均衡算法限流算法
静态策略随机、哈希、轮询固定窗口、滑动窗口、漏桶、令牌桶
动态策略最少连接数、最少活跃请求数BBR

7.2 静态策略的缺陷

静态算法依赖于程序员的个人经验。例如在限流中,有三种方式确定限流的阈值:

  1. 基于压测结果——通过压力测试找到系统的承受能力
  2. 基于可观测性数据——根据线上监控数据调整
  3. 基于源码分析——分析代码逻辑推算

对于负载均衡,有一些基本假设:

  • 所有的请求都消耗一样的资源
  • 大量的请求

如果按照 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[客户端使用指标
执行负载均衡]
  1. 可观测性平台从所有的节点采集指标
  2. 客户端拉取指标(或者可观测性平台推送指标)
  3. 客户端使用这些指标来执行负载均衡

开发者可以根据指标和业务特征来设计负载均衡算法。

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)的核心概念和代码实现,最后探讨了如何将可观测性数据与服务治理结合起来。

关键知识点回顾:

  1. AOP 是可用性和可观测性的基础——通过拦截器在不侵入业务代码的前提下增加治理逻辑
  2. 熔断、限流、降级本质相同——都是故障检测 + 故障处理 + 故障恢复,区别在于处理策略
  3. 限流算法中令牌桶和漏桶效果一样,滑动窗口比固定窗口更平滑
  4. 熔断器是一个三状态状态机:Closed -> Open -> Half-Open
  5. 集群限流通过 Redis + Lua 脚本实现,保证集群维度的原子性
  6. 可观测性三支柱:Metrics(指标)、Tracing(链路追踪)、Logging(日志)
  7. 链路追踪通过 TraceID 将多个服务的调用串联起来,核心是 Extract 和 Inject
  8. 基于可观测性的服务治理是将监控数据与负载均衡、限流等治理策略结合
  9. 指标的时间敏感性是动态治理的核心挑战,综合方案(响应附带 + 平台采集)是最优解
  10. 服务端治理比客户端治理简单,因为服务端拥有自己的全部实时信息

自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 熔断、限流、降级三者本质是否不同?它们之间最核心的区别是什么?
  2. 滑动窗口相比固定窗口解决了什么问题?为什么固定窗口在窗口边界可能放行近 2 倍流量?
  3. 令牌桶和漏桶"效果一样"指的是什么?令牌桶相对漏桶多出的能力是什么?
  4. 熔断器三个状态(Closed / Open / Half-Open)之间如何转换?引入 Half-Open 半开状态是为了解决什么?
  5. 可观测性三支柱分别回答什么问题?TraceID 和 SpanID(含 ParentSpanID)各自的职责是什么?

动手练习(建议真做一遍)

  1. 把本文的令牌桶 Allow() 改造成"阻塞等待"版本(拿不到令牌就 time.Sleep 一小会儿重试,直到拿到或整体超时),对比"直接拒绝"与"阻塞排队"两种策略在突发流量下的表现。
  2. 把 6.1 的 MetricsInterceptor 真正接入 Prometheus(用官方 client_golang 计数器与直方图),用 grpc_request_duration_seconds 看到 QPS 与 P99 曲线。
  3. 在 5.3 Redis 集群限流的 RedisFixedWindow 上,故意把 INCREXPIRE 拆成两条非原子命令,写一个并发脚本复现"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 上,构建可落地的日志、追踪与监控体系。

About Me

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

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

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

目标

学AI,加油!加油!