前置知识: Go、Go、Go、Go

Go 与限流

19 min高级

Go限流与熔断详解:令牌桶、漏桶、滑动窗口、分布式限流、熔断器、自适应限流与工程实践

前置知识

学习目标

  • 掌握「1. 历史动机与背景」的核心机制、典型用法与常见陷阱
  • 掌握「2. 形式化定义」的核心机制、典型用法与常见陷阱
  • 掌握「3. 理论推导」的核心机制、典型用法与常见陷阱
  • 掌握「4. 代码示例」的核心机制、典型用法与常见陷阱
  • 掌握「5. 算法对比与选型」的核心机制、典型用法与常见陷阱

1. 历史动机与背景

1.1 限流算法的起源

限流(Rate Limiting)最早源于电信网络的流量控制(Traffic Policing)。20 世纪 60 年代,AT&T Bell Labs 在分组交换网络研究中提出了最早的令牌桶算法:

  • 1965:IBM 研究员 J. R. W. Smith 提出”bucket algorithm”概念,用于控制数据传输速率。
  • 1989:ITU-T(国际电信联盟)在 I.371 标准中正式定义了”Token Bucket”与”Leaky Bucket”两种算法。
  • 1995:Cisco 路由器首次实现基于令牌桶的流量监管(Committed Access Rate, CAR)。
  • 1999:Google 在搜索后端引入限流机制,防止爬虫拖垮核心服务。

1.2 熔断器的起源

熔断器(Circuit Breaker)概念源于电气工程,由 Michael Nygard 在 2007 年的著作 Release It! 中首次引入软件工程领域:

  • 2007:Nygard 提出 Circuit Breaker 模式,保护下游服务调用。
  • 2012:Netflix 开源 Hystrix,将熔断、隔离、降级整合为完整容错框架。
  • 2015:Go 社区出现 sony/gobreaker,提供轻量级熔断实现。
  • 2017:Hystrix 停止维护,社区转向 Resilience4j(Java)与各语言的轻量实现。
  • 2019:Netflix 发布 Concurrency Limits,提出自适应限流思想。

1.3 Go 限流生态演进

年份工具/库说明
2015golang.org/x/time/rateGo 官方扩展库提供令牌桶实现
2015github.com/sony/gobreakerSony 开源的 Go 熔断器实现
2016github.com/afex/hystrix-goNetflix Hystrix 的 Go 移植版
2017github.com/ulule/limiter支持多种后端(Redis、Memcached)的限流库
2018github.com/throttled/throttled支持 Redis 的分布式限流库
2019github.com/go-redis/redis_rate基于 Redis 的令牌桶限流
2020github.com/cep21/circuitGo 熔断器库(受 Hystrix 启发)
2022github.com/sony/gobreaker/v2gobreaker v2 发布,支持泛型
2023github.com/sourcegraph/conc提供 pool、waitgroup 等并发原语

1.4 为什么需要限流与熔断

分布式系统中的稳定性问题:

  1. 突发流量(Flash Crowd):秒杀、热点新闻、DDoS 攻击可能导致请求量瞬间激增 10-100 倍。
  2. 雪崩效应(Cascading Failure):下游服务慢响应导致上游连接堆积,拖垮整个调用链。
  3. 资源耗尽(Resource Exhaustion):连接池、线程池、内存等有限资源被耗尽,拒绝正常请求。
  4. 级联故障(Cascading Failure):单个服务故障通过同步调用传播至整个系统。

限流与熔断的分工:

  • 限流:控制进入系统的请求量,保护自身不被打垮。
  • 熔断:控制调用下游的请求量,保护自身不被下游拖垮。
  • 隔离:限制不同业务间资源争抢(信号量、Bulkhead)。
  • 降级:在故障时返回兜底响应,保证核心链路可用。

1.5 与其他语言生态对比

  • Java:Resilience4j(函数式 + 装饰器模式)、Hystrix(已停维护)、Sentinel(阿里,支持流控、熔断、热点限流)。
  • Rust:tower 框架提供限流与熔断中间件,governor 实现令牌桶。
  • Python:ratelimit 装饰器、aiometer 异步限流。
  • Node.js:express-rate-limit、bottleneck。

Go 的优势:

  1. 标准库官方维护:golang.org/x/time/rate 由 Go 团队维护,质量有保障。
  2. 轻量无依赖:gobreaker 仅数百行代码,无第三方依赖。
  3. goroutine 友好:rate.Limiter 基于时间戳而非后台 goroutine,无锁竞争。

2. 形式化定义

2.1 令牌桶算法的形式化

令牌桶(Token Bucket)由两个参数描述:

  • rr:令牌生成速率(tokens per second),即每秒补充的令牌数。
  • bb:桶容量(burst),桶能容纳的最大令牌数。

设 T(t)T(t) 为时刻 tt 桶中的令牌数,则令牌数随时间变化:

T(t)=min⁡(b, T(t0)+r⋅(t−t0))T(t) = \min\left(b, \, T(t_0) + r \cdot (t - t_0)\right)

其中 t0t_0 是上次更新时间。请求消耗令牌:

Allow(n)={true若 T(t)≥nfalse否则\text{Allow}(n) = \begin{cases} \text{true} & \text{若 } T(t) \geq n \\ \text{false} & \text{否则} \end{cases}

若允许则 T(t)←T(t)−nT(t) \leftarrow T(t) - n。

2.2 漏桶算法的形式化

漏桶(Leaky Bucket)以固定速率 rr “漏出”(处理)请求,桶容量 bb 限制排队请求数:

Q(t)=min⁡(b, Q(t0)+Arrivals(t0,t)−r⋅(t−t0))Q(t) = \min\left(b, \, Q(t_0) + \text{Arrivals}(t_0, t) - r \cdot (t - t_0)\right)

其中 Q(t)Q(t) 是队列长度,Arrivals(t0,t)\text{Arrivals}(t_0, t) 是 (t0,t)(t_0, t) 内到达的请求数。

令牌桶与漏桶的核心差异:

  • 令牌桶允许瞬间突发 bb 个请求(桶满时),平均速率受 rr 限制。
  • 漏桶强制匀速 rr,超出 bb 的请求直接丢弃。

2.3 滑动窗口算法的形式化

设窗口长度为 WW,当前时刻 tt,滑动窗口内请求数:

N(t)=∣{req∣t−W≤time(req)≤t}∣N(t) = \left| \{ \text{req} \mid t - W \leq \text{time}(\text{req}) \leq t \} \right|

限流决策:

Allow(t)={true若 N(t)<Lfalse否则\text{Allow}(t) = \begin{cases} \text{true} & \text{若 } N(t) < L \\ \text{false} & \text{否则} \end{cases}

其中 LL 是窗口内允许的最大请求数。

滑动窗口的精度优化:将窗口分为 kk 个小格子(cell),每个格子记录一个时间段的请求数:

N(t)=∑i=0k−1celli(t−i⋅W/k)N(t) = \sum_{i=0}^{k-1} \text{cell}_i(t - i \cdot W/k)

复杂度从 O(N)O(N) 降至 O(k)O(k),kk 通常取 5-10。

2.4 熔断器状态机的形式化

熔断器是一个三状态有限自动机:

State={Closed,Open,HalfOpen}\text{State} = \{ \text{Closed}, \text{Open}, \text{HalfOpen} \}

状态转移:

Closed→failures≥thresholdOpenOpen→t>TimeoutHalfOpenHalfOpen→successes≥successThresholdClosedHalfOpen→failureOpen\begin{aligned} \text{Closed} &\xrightarrow{\text{failures} \geq \text{threshold}} \text{Open} \\ \text{Open} &\xrightarrow{t > \text{Timeout}} \text{HalfOpen} \\ \text{HalfOpen} &\xrightarrow{\text{successes} \geq \text{successThreshold}} \text{Closed} \\ \text{HalfOpen} &\xrightarrow{\text{failure}} \text{Open} \end{aligned}

在 Open 状态下,所有请求直接失败(快速失败,fail-fast);在 HalfOpen 状态下,允许有限个测试请求通过。

2.5 限流决策的代价模型

限流的代价包括:

Climit=Cdecision+Cstate+CstorageC_{\text{limit}} = C_{\text{decision}} + C_{\text{state}} + C_{\text{storage}}
  • CdecisionC_{\text{decision}}:算法决策开销。令牌桶 O(1)O(1),滑动窗口 O(k)O(k)。
  • CstateC_{\text{state}}:状态维护开销。单机内存 O(1)O(1),分布式需同步。
  • CstorageC_{\text{storage}}:状态存储开销。单机 sync.Map O(N)O(N),Redis O(N)O(N)。

对 rate.Limiter:Climit≈O(1)C_{\text{limit}} \approx O(1)(原子操作 + 时间戳计算)。


3. 理论推导

3.1 令牌桶的无锁实现

golang.org/x/time/rate 的核心实现基于”惰性补充”(lazy refill):

  1. 记录上次更新时间:last time.Time。
  2. 请求时计算当前可用令牌数:
tokens(t)=min⁡(b, tokens(t0)+r⋅(t−t0))\text{tokens}(t) = \min\left(b, \, \text{tokens}(t_0) + r \cdot (t - t_0)\right)
  1. 判断与消耗:若 tokens(t)≥n\text{tokens}(t) \geq n 则消耗 nn 个令牌,更新 tokens\text{tokens} 与 t0t_0。

这种实现避免了后台 goroutine 不断补充令牌,仅需原子操作即可。伪代码:

// Allow 判断是否允许消耗 n 个令牌
func (lim *Limiter) AllowN(now time.Time, n int) bool {
    lim.mu.Lock()
    defer lim.mu.Unlock()

    // 计算从上次到现在新增的令牌
    elapsed := now.Sub(lim.last)
    delta := int(elapsed.Seconds() * float64(lim.rate))
    tokens := min(lim.burst, lim.tokens + delta)

    if tokens >= n {
        lim.tokens = tokens - n
        lim.last = now
        return true
    }
    return false
}

注意:实际实现使用浮点数表示令牌数,以支持小于 1 的速率(如 0.5 token/s)。

3.2 Reserve 方法的延迟计算

Reserve 方法不立即拒绝请求,而是计算需要等待多久才能获取令牌:

Delay(t,n)=max⁡(0, n−tokens(t)r)\text{Delay}(t, n) = \max\left(0, \, \frac{n - \text{tokens}(t)}{r}\right)

若 n>bn > b(请求令牌数超过桶容量),Delay 无穷大,Reserve 返回 OK=false。

使用 time.Sleep(reservation.Delay()) 可精确等待,实现”匀速 + 排队”语义。

3.3 滑动窗口的 Redis 实现

使用 Redis ZSET 存储请求时间戳:

  1. 移除窗口外记录:ZREMRANGEBYSCORE key 0 (now - window)。
  2. 添加当前请求:ZADD key now member。
  3. 统计窗口内数量:ZCARD key。
  4. 设置过期:EXPIRE key window(避免内存泄漏)。

复杂度:O(log⁡N+N)O(\log N + N),NN 为窗口内请求数。

优化:使用 Lua 脚本保证原子性:

-- KEYS[1]: 限流 key
-- ARGV[1]: 当前时间戳
-- ARGV[2]: 窗口大小(秒)
-- ARGV[3]: 限制数量
local key = KEYS[1]
local now = tonumber(ARGV[1])
local window = tonumber(ARGV[2])
local limit = tonumber(ARGV[3])

redis.call('ZREMRANGEBYSCORE', key, 0, now - window * 1000)
local count = redis.call('ZCARD', key)
if count < limit then
    redis.call('ZADD', key, now, now)
    redis.call('PEXPIRE', key, window * 1000)
    return 1
else
    return 0
end

3.4 令牌桶的分布式实现

Redis 令牌桶需保证原子性。常用方案:

  1. Redis + Lua 脚本:在 Redis 内计算令牌数,原子操作。
  2. Redis Cell 模块:Redis 官方的限流模块(CL.THROTTLE 命令)。

Lua 脚本核心逻辑:

-- 令牌桶:维护上次更新时间与剩余令牌
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])

local data = redis.call('HMGET', key, 'tokens', 'last')
local tokens = tonumber(data[1]) or capacity
local last = tonumber(data[2]) or now

local delta = max(0, now - last) * rate
tokens = min(capacity, tokens + delta)

local allowed = tokens >= 1
if allowed then
    tokens = tokens - 1
end

redis.call('HMSET', key, 'tokens', tokens, 'last', now)
redis.call('EXPIRE', key, 60)

return allowed

3.5 熔断器的状态转移推导

熔断器的核心是”失败统计”与”状态转移”。

设时间窗口 WW 内的总请求数为 NN,失败数为 FF,则失败率:

failRate=FN\text{failRate} = \frac{F}{N}

熔断触发条件:

trip=(N≥minRequests)∧(failRate≥threshold∨F≥consecutiveFailures)\text{trip} = (N \geq \text{minRequests}) \land \left( \text{failRate} \geq \text{threshold} \lor F \geq \text{consecutiveFailures} \right)

其中 minRequests 避免少量请求时的误判(如 3 次请求失败 2 次,失败率 67%,但样本不足)。

恢复策略:

  • 定时恢复:Open 状态持续 Timeout 时间后转入 HalfOpen。
  • 探测恢复:HalfOpen 状态下允许 MaxRequests 个请求通过,若全部成功则恢复 Closed,任一失败则回到 Open。

3.6 复杂度分析

操作时间复杂度备注
rate.Limiter.AllowO(1)O(1)原子操作 + 时间戳计算
rate.Limiter.WaitO(delay)O(\text{delay})阻塞等待
滑动窗口(内存)O(k)O(k)kk 个格子
滑动窗口(Redis ZSET)O(log⁡N)O(\log N)ZSET 操作
熔断器状态更新O(1)O(1)计数器更新
按客户端限流查找O(1)O(1) 均摊sync.Map

4. 代码示例

4.1 基础:rate.Limiter 基本用法

// Package main 演示 golang.org/x/time/rate 的基本用法
package main

import (
    "context"
    "fmt"
    "time"

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

func main() {
    // 创建限流器:每秒 10 个令牌,桶容量 5(允许 5 个突发)
    limiter := rate.NewLimiter(10, 5)

    // 方式1:Allow - 立即返回是否允许
    for i := 0; i < 10; i++ {
        if limiter.Allow() {
            fmt.Printf("请求 %d 通过\n", i+1)
        } else {
            fmt.Printf("请求 %d 被限流\n", i+1)
        }
    }

    // 方式2:Wait - 阻塞等待直到获取令牌
    ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
    defer cancel()
    if err := limiter.Wait(ctx); err != nil {
        fmt.Println("等待超时或被取消:", err)
    }

    // 方式3:Reserve - 预留令牌,返回需要等待的时间
    reservation := limiter.Reserve()
    if !reservation.OK() {
        fmt.Println("令牌数超过桶容量")
        return
    }
    delay := reservation.Delay()
    fmt.Printf("需要等待 %v\n", delay)
    time.Sleep(delay)

    // 动态调整速率
    limiter.SetLimit(20) // 调整为每秒 20 个
    limiter.SetBurst(10) // 调整桶容量为 10
}

4.2 HTTP 中间件限流

// Package middleware 提供 HTTP 限流中间件
package middleware

import (
    "net/http"

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

// RateLimitMiddleware 全局限流中间件
// 参数:
//   - limiter: 限流器实例
// 返回:HTTP 中间件函数
func RateLimitMiddleware(limiter *rate.Limiter) func(http.Handler) http.Handler {
    return func(next http.Handler) http.Handler {
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
            if !limiter.Allow() {
                // 设置 Retry-After 头,告知客户端重试等待时间
                w.Header().Set("Retry-After", "1")
                http.Error(w, "请求过于频繁,请稍后再试", http.StatusTooManyRequests)
                return
            }
            next.ServeHTTP(w, r)
        })
    }
}

// 使用示例
// limiter := rate.NewLimiter(100, 10) // 每秒 100 个,突发 10
// handler := RateLimitMiddleware(limiter)(mux)
// http.ListenAndServe(":8080", handler)

4.3 按客户端限流(IP 维度)

// Package middleware 提供按 IP 限流中间件
package middleware

import (
    "net"
    "net/http"
    "sync"

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

// IPRateLimiter 按 IP 限流器
type IPRateLimiter struct {
    limiters sync.Map // map[string]*rate.Limiter
    rate     rate.Limit
    burst    int
}

// NewIPRateLimiter 创建按 IP 限流器
// 参数:
//   - r: 每秒令牌数
//   - burst: 桶容量
func NewIPRateLimiter(r rate.Limit, burst int) *IPRateLimiter {
    return &IPRateLimiter{
        rate:  r,
        burst: burst,
    }
}

// GetLimiter 获取或创建指定 IP 的限流器
func (l *IPRateLimiter) GetLimiter(ip string) *rate.Limiter {
    if limiter, ok := l.limiters.Load(ip); ok {
        return limiter.(*rate.Limiter)
    }
    limiter := rate.NewLimiter(l.rate, l.burst)
    l.limiters.Store(ip, limiter)
    return limiter
}

// Middleware 返回 HTTP 中间件
func (l *IPRateLimiter) Middleware(next http.Handler) http.Handler {
    return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
        // 提取客户端 IP(考虑反向代理)
        ip, _, err := net.SplitHostPort(r.RemoteAddr)
        if err != nil {
            ip = r.RemoteAddr
        }
        if forwarded := r.Header.Get("X-Forwarded-For"); forwarded != "" {
            ip = forwarded
        }

        if !l.GetLimiter(ip).Allow() {
            w.Header().Set("Retry-After", "1")
            http.Error(w, "请求过于频繁", http.StatusTooManyRequests)
            return
        }
        next.ServeHTTP(w, r)
    })
}

4.4 按 IP 限流 + LRU 淘汰(防内存泄漏)

// Package middleware 提供带 LRU 淘汰的按 IP 限流器
package middleware

import (
    "container/list"
    "sync"
    "time"

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

// ipEntry 限流器条目
type ipEntry struct {
    ip      string
    limiter *rate.Limiter
    lastUse time.Time
}

// LRUIPRateLimiter 带 LRU 淘汰的按 IP 限流器
type LRUIPRateLimiter struct {
    mu       sync.Mutex
    entries  map[string]*list.Element
    lru      *list.List
    rate     rate.Limit
    burst    int
    maxSize  int           // 最大 IP 数量
    ttl      time.Duration  // 不活跃超时
}

// NewLRUIPRateLimiter 创建带 LRU 的限流器
func NewLRUIPRateLimiter(r rate.Limit, burst, maxSize int, ttl time.Duration) *LRUIPRateLimiter {
    l := &LRUIPRateLimiter{
        entries: make(map[string]*list.Element),
        lru:     list.New(),
        rate:    r,
        burst:   burst,
        maxSize: maxSize,
        ttl:     ttl,
    }
    // 启动后台清理 goroutine
    go l.cleanup()
    return l
}

// GetLimiter 获取限流器,同时更新 LRU
func (l *LRUIPRateLimiter) GetLimiter(ip string) *rate.Limiter {
    l.mu.Lock()
    defer l.mu.Unlock()

    if elem, ok := l.entries[ip]; ok {
        l.lru.MoveToFront(elem)
        entry := elem.Value.(*ipEntry)
        entry.lastUse = time.Now()
        return entry.limiter
    }

    // 创建新条目
    entry := &ipEntry{
        ip:      ip,
        limiter: rate.NewLimiter(l.rate, l.burst),
        lastUse: time.Now(),
    }
    elem := l.lru.PushFront(entry)
    l.entries[ip] = elem

    // 超出容量则淘汰最久未使用
    if l.lru.Len() > l.maxSize {
        oldest := l.lru.Back()
        if oldest != nil {
            entry := oldest.Value.(*ipEntry)
            delete(l.entries, entry.ip)
            l.lru.Remove(oldest)
        }
    }
    return entry.limiter
}

// cleanup 定期清理不活跃的 IP 限流器
func (l *LRUIPRateLimiter) cleanup() {
    ticker := time.NewTicker(time.Minute)
    defer ticker.Stop()
    for range ticker.C {
        l.mu.Lock()
        now := time.Now()
        var toRemove []string
        for ip, elem := range l.entries {
            entry := elem.Value.(*ipEntry)
            if now.Sub(entry.lastUse) > l.ttl {
                toRemove = append(toRemove, ip)
            }
        }
        for _, ip := range toRemove {
            if elem, ok := l.entries[ip]; ok {
                l.lru.Remove(elem)
                delete(l.entries, ip)
            }
        }
        l.mu.Unlock()
    }
}

4.5 Redis 滑动窗口限流

// Package ratelimit 提供 Redis 分布式限流
package ratelimit

import (
    "context"
    "errors"
    "time"

    "github.com/redis/go-redis/v9"
)

// SlidingWindowLimiter 滑动窗口限流器
type SlidingWindowLimiter struct {
    rdb    *redis.Client
    key    string        // 限流 key(如 "rate:user:123")
    limit  int           // 窗口内最大请求数
    window time.Duration // 窗口大小
    script *redis.Script // 预编译 Lua 脚本
}

// NewSlidingWindowLimiter 创建滑动窗口限流器
func NewSlidingWindowLimiter(rdb *redis.Client, key string, limit int, window time.Duration) *SlidingWindowLimiter {
    return &SlidingWindowLimiter{
        rdb:    rdb,
        key:    key,
        limit:  limit,
        window: window,
        script: redis.NewScript(slidingWindowScript),
    }
}

// 滑动窗口 Lua 脚本(保证原子性)
const slidingWindowScript = `
local key = KEYS[1]
local now = tonumber(ARGV[1])
local window = tonumber(ARGV[2])
local limit = tonumber(ARGV[3])

-- 移除窗口外的记录
redis.call('ZREMRANGEBYSCORE', key, 0, now - window)

-- 统计当前窗口内请求数
local count = redis.call('ZCARD', key)

if count < limit then
    -- 添加当前请求
    redis.call('ZADD', key, now, now)
    -- 设置过期时间,避免内存泄漏
    redis.call('PEXPIRE', key, window)
    return 1
else
    return 0
end
`

// Allow 判断是否允许通过
func (l *SlidingWindowLimiter) Allow(ctx context.Context) (bool, error) {
    now := time.Now().UnixMilli()
    windowMs := l.window.Milliseconds()
    result, err := l.script.Run(ctx, l.rdb, []string{l.key}, now, windowMs, l.limit).Int()
    if err != nil {
        return false, err
    }
    return result == 1, nil
}

// 使用示例
func ExampleUsage() {
    rdb := redis.NewClient(&redis.Options{Addr: "localhost:6379"})
    limiter := NewSlidingWindowLimiter(rdb, "rate:user:123", 100, time.Minute)

    ctx := context.Background()
    allowed, err := limiter.Allow(ctx)
    if err != nil {
        // 处理错误(如 Redis 不可用)
        // 实践中应降级为单机限流或直接拒绝
        _ = errors.New("rate limit check failed")
    }
    _ = allowed
}

4.6 Redis 令牌桶限流

// Package ratelimit 提供 Redis 令牌桶限流
package ratelimit

import (
    "context"
    "time"

    "github.com/redis/go-redis/v9"
)

// TokenBucketLimiter Redis 令牌桶限流器
type TokenBucketLimiter struct {
    rdb      *redis.Client
    key      string
    capacity int           // 桶容量
    rate     float64       // 令牌生成速率(tokens/s)
    script   *redis.Script
}

// 令牌桶 Lua 脚本
const tokenBucketScript = `
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local requested = tonumber(ARGV[4])

-- 获取当前状态
local data = redis.call('HMGET', key, 'tokens', 'last')
local tokens = tonumber(data[1]) or capacity
local last = tonumber(data[2]) or now

-- 计算补充的令牌
local elapsed = math.max(0, now - last)
local refill = elapsed * rate / 1000  -- now 单位为毫秒
tokens = math.min(capacity, tokens + refill)

-- 判断与消耗
local allowed = 0
if tokens >= requested then
    tokens = tokens - requested
    allowed = 1
end

-- 更新状态
redis.call('HMSET', key, 'tokens', tokens, 'last', now)
redis.call('EXPIRE', key, 60)

return allowed
`

// NewTokenBucketLimiter 创建 Redis 令牌桶限流器
func NewTokenBucketLimiter(rdb *redis.Client, key string, capacity int, rate float64) *TokenBucketLimiter {
    return &TokenBucketLimiter{
        rdb:      rdb,
        key:      key,
        capacity: capacity,
        rate:      rate,
        script:   redis.NewScript(tokenBucketScript),
    }
}

// Allow 判断是否允许通过
func (l *TokenBucketLimiter) Allow(ctx context.Context) (bool, error) {
    now := time.Now().UnixMilli()
    result, err := l.script.Run(ctx, l.rdb, []string{l.key}, l.capacity, l.rate, now, 1).Int()
    if err != nil {
        return false, err
    }
    return result == 1, nil
}

4.7 熔断器(gobreaker)

// Package breaker 提供 gobreaker 熔断器封装
package breaker

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

    "github.com/sony/gobreaker"
)

// CircuitBreaker 熔断器封装
type CircuitBreaker struct {
    cb *gobreaker.CircuitBreaker
}

// Config 熔断器配置
type Config struct {
    Name          string          // 熔断器名称
    MaxRequests   uint32          // 半开状态下允许的测试请求数
    Interval      time.Duration   // Closed 状态下的统计周期
    Timeout       time.Duration   // Open 状态持续多久后进入 HalfOpen
    FailThreshold uint32          // 连续失败多少次触发熔断
    FailRate      float64         // 失败率阈值(0-1)
    MinRequests   uint32          // 触发失败率判断的最小请求数
}

// NewCircuitBreaker 创建熔断器
func NewCircuitBreaker(cfg Config) *CircuitBreaker {
    settings := gobreaker.Settings{
        Name:        cfg.Name,
        MaxRequests: cfg.MaxRequests,
        Interval:    cfg.Interval,
        Timeout:     cfg.Timeout,
        ReadyToTrip: func(counts gobreaker.Counts) bool {
            // 连续失败数触发
            if counts.ConsecutiveFailures > cfg.FailThreshold {
                return true
            }
            // 失败率触发(需达到最小请求数)
            if counts.Requests >= uint32(cfg.MinRequests) {
                failRate := float64(counts.TotalFailures) / float64(counts.Requests)
                if failRate >= cfg.FailRate {
                    return true
                }
            }
            return false
        },
        OnStateChange: func(name string, from, to gobreaker.State) {
            // 状态变更回调(可用于监控、告警)
            fmt.Printf("[Breaker %s] %s -> %s\n", name, from, to)
        },
    }
    return &CircuitBreaker{cb: gobreaker.NewCircuitBreaker(settings)}
}

// ErrCircuitOpen 熔断器打开错误
var ErrCircuitOpen = errors.New("circuit breaker is open")

// Execute 执行请求,自动熔断
func (b *CircuitBreaker) Execute(ctx context.Context, fn func(ctx context.Context) (interface{}, error)) (interface{}, error) {
    result, err := b.cb.Execute(func() (interface{}, error) {
        return fn(ctx)
    })
    if err == gobreaker.ErrOpenCircuit {
        return nil, ErrCircuitOpen
    }
    if errors.Is(err, gobreaker.ErrTooManyRequests) {
        return nil, errors.New("half-open: too many requests")
    }
    return result, err
}

// State 返回当前状态
func (b *CircuitBreaker) State() gobreaker.State {
    return b.cb.State()
}

// Counts 返回统计信息
func (b *CircuitBreaker) Counts() gobreaker.Counts {
    return b.cb.Counts()
}

// 使用示例
type UserService struct {
    breaker *CircuitBreaker
    client  *http.Client
}

func (s *UserService) GetUser(ctx context.Context, id string) (*User, error) {
    result, err := s.breaker.Execute(ctx, func(ctx context.Context) (interface{}, error) {
        // 调用下游服务
        req, _ := http.NewRequestWithContext(ctx, "GET", "http://user-service/users/"+id, nil)
        resp, err := s.client.Do(req)
        if err != nil {
            return nil, err
        }
        defer resp.Body.Close()
        if resp.StatusCode >= 500 {
            return nil, fmt.Errorf("server error: %d", resp.StatusCode)
        }
        if resp.StatusCode == 404 {
            return nil, ErrNotFound
        }
        var user User
        if err := json.NewDecoder(resp.Body).Decode(&user); err != nil {
            return nil, err
        }
        return &user, nil
    })
    if err != nil {
        return nil, err
    }
    return result.(*User), nil
}

4.8 并发限制(信号量)

// Package middleware 提供并发数限制中间件
package middleware

import (
    "net/http"

    "golang.org/x/sync/semaphore"
)

// ConcurrencyLimitMiddleware 并发数限制中间件
// 参数:
//   - max: 最大并发数
func ConcurrencyLimitMiddleware(max int64) func(http.Handler) http.Handler {
    sem := semaphore.NewWeighted(max)
    return func(next http.Handler) http.Handler {
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
            // 尝试获取信号量(非阻塞)
            if !sem.TryAcquire(1) {
                http.Error(w, "服务器繁忙,请稍后再试", http.StatusServiceUnavailable)
                return
            }
            defer sem.Release(1)
            next.ServeHTTP(w, r)
        })
    }
}

// ConcurrencyLimitWithWait 带等待的并发数限制
func ConcurrencyLimitWithWait(max int64) func(http.Handler) http.Handler {
    sem := semaphore.NewWeighted(max)
    return func(next http.Handler) http.Handler {
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
            // 阻塞等待信号量(支持 Context 取消)
            if err := sem.Acquire(r.Context(), 1); err != nil {
                http.Error(w, "请求超时", http.StatusServiceUnavailable)
                return
            }
            defer sem.Release(1)
            next.ServeHTTP(w, r)
        })
    }
}

4.9 自适应限流

// Package ratelimit 提供基于 CPU 与延迟的自适应限流
package ratelimit

import (
    "math"
    "runtime"
    "sync"
    "sync/atomic"
    "time"

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

// AdaptiveLimiter 自适应限流器
// 根据 CPU 使用率与请求延迟动态调整限流阈值
type AdaptiveLimiter struct {
    mu          sync.Mutex
    limiter     *rate.Limiter
    minRate     rate.Limit // 最小速率
    maxRate     rate.Limit // 最大速率
    currentRate rate.Limit

    // 指标采集
    cpuUsage       float64 // 0-1
    p99Latency     time.Duration
    latencyThreshold time.Duration

    // 窗口内统计
    requests  int64
    failures  int64
    lastReset time.Time
}

// NewAdaptiveLimiter 创建自适应限流器
func NewAdaptiveLimiter(initialRate, minRate, maxRate rate.Limit, latencyThreshold time.Duration) *AdaptiveLimiter {
    return &AdaptiveLimiter{
        limiter:          rate.NewLimiter(initialRate, int(initialRate)),
        minRate:          minRate,
        maxRate:          maxRate,
        currentRate:      initialRate,
        latencyThreshold: latencyThreshold,
        lastReset:        time.Now(),
    }
}

// Allow 判断是否允许通过
func (l *AdaptiveLimiter) Allow() bool {
    return l.limiter.Allow()
}

// Record 记录请求结果
func (l *AdaptiveLimiter) Record(latency time.Duration, success bool) {
    atomic.AddInt64(&l.requests, 1)
    if !success {
        atomic.AddInt64(&l.failures, 1)
    }

    // 更新 P99 延迟(简化版,实际需用滑动窗口或 t-digest)
    l.mu.Lock()
    if latency > l.p99Latency {
        l.p99Latency = latency
    }
    l.mu.Unlock()
}

// Adjust 根据系统负载动态调整速率
// CPU 使用率高或延迟超阈值则降低速率,反之提高
func (l *AdaptiveLimiter) Adjust() {
    l.mu.Lock()
    defer l.mu.Unlock()

    now := time.Now()
    if now.Sub(l.lastReset) < time.Second {
        return
    }

    // 计算窗口内失败率
    reqs := atomic.LoadInt64(&l.requests)
    fails := atomic.LoadInt64(&l.failures)
    failRate := 0.0
    if reqs > 0 {
        failRate = float64(fails) / float64(reqs)
    }

    // AIMD(Additive Increase, Multiplicative Decrease)策略
    switch {
    case l.cpuUsage > 0.8 || failRate > 0.1 || l.p99Latency > l.latencyThreshold:
        // 过载:乘性降低(降低 20%)
        newRate := rate.Limit(float64(l.currentRate) * 0.8)
        if newRate < l.minRate {
            newRate = l.minRate
        }
        l.setRate(newRate)
    case l.cpuUsage < 0.5 && failRate < 0.01:
        // 轻载:加性增加(提升 10%)
        newRate := rate.Limit(float64(l.currentRate) * 1.1)
        if newRate > l.maxRate {
            newRate = l.maxRate
        }
        l.setRate(newRate)
    }

    // 重置统计
    atomic.StoreInt64(&l.requests, 0)
    atomic.StoreInt64(&l.failures, 0)
    l.lastReset = now
    l.p99Latency = 0
}

// setRate 设置新速率
func (l *AdaptiveLimiter) setRate(r rate.Limit) {
    if r != l.currentRate {
        l.currentRate = r
        l.limiter.SetLimit(r)
        l.limiter.SetBurst(int(math.Max(1, float64(r))))
    }
}

// SetCPUUsage 更新 CPU 使用率(由外部采集)
func (l *AdaptiveLimiter) SetCPUUsage(usage float64) {
    l.mu.Lock()
    l.cpuUsage = usage
    l.mu.Unlock()
}

// CPU 使用率采集 goroutine
func collectCPUUsage(limiter *AdaptiveLimiter, interval time.Duration) {
    var lastCPU time.Time
    var lastStats *runtime.MemStats
    _ = lastStats

    ticker := time.NewTicker(interval)
    defer ticker.Stop()

    for range ticker.C {
        // 简化版:使用 runtime 估算 CPU 使用率
        // 生产中应读取 /proc/stat(Linux)或用 gopsutil 库
        now := time.Now()
        _ = now.Sub(lastCPU)
        lastCPU = now

        // 占位:实际应采集真实 CPU 使用率
        // limiter.SetCPUUsage(cpuUsage)
    }
}

4.10 完整限流中间件组合

// Package middleware 提供组合限流中间件
package middleware

import (
    "net/http"
    "time"

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

// CompositeLimitConfig 组合限流配置
type CompositeLimitConfig struct {
    GlobalRate     rate.Limit // 全局速率
    GlobalBurst    int        // 全局突发
    UserRate       rate.Limit // 单用户速率
    UserBurst      int        // 单用户突发
    MaxConcurrency int64      // 最大并发数
    RetryAfter     int        // Retry-After 头值(秒)
}

// NewCompositeMiddleware 创建组合限流中间件
// 包含:全局限流 + 按用户限流 + 并发限制
func NewCompositeMiddleware(cfg CompositeLimitConfig) func(http.Handler) http.Handler {
    globalLimiter := rate.NewLimiter(cfg.GlobalRate, cfg.GlobalBurst)
    userLimiter := NewLRUIPRateLimiter(cfg.UserRate, cfg.UserBurst, 100000, 10*time.Minute)
    concurrencySem := semaphore.NewWeighted(cfg.MaxConcurrency)

    return func(next http.Handler) http.Handler {
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
            // 1. 全局限流
            if !globalLimiter.Allow() {
                respondRateLimited(w, cfg.RetryAfter, "系统繁忙,请稍后再试")
                return
            }

            // 2. 用户限流
            ip := extractIP(r)
            if !userLimiter.GetLimiter(ip).Allow() {
                respondRateLimited(w, cfg.RetryAfter, "请求过于频繁")
                return
            }

            // 3. 并发限制(非阻塞)
            if !concurrencySem.TryAcquire(1) {
                http.Error(w, "服务器繁忙", http.StatusServiceUnavailable)
                return
            }
            defer concurrencySem.Release(1)

            next.ServeHTTP(w, r)
        })
    }
}

// respondRateLimited 返回标准限流响应
func respondRateLimited(w http.ResponseWriter, retryAfter int, message string) {
    w.Header().Set("Retry-After", fmt.Sprintf("%d", retryAfter))
    w.Header().Set("Content-Type", "application/json")
    w.WriteHeader(http.StatusTooManyRequests)
    json.NewEncoder(w).Encode(map[string]interface{}{
        "error":       "too_many_requests",
        "message":     message,
        "retry_after": retryAfter,
    })
}

// extractIP 提取客户端 IP
func extractIP(r *http.Request) string {
    if forwarded := r.Header.Get("X-Forwarded-For"); forwarded != "" {
        return forwarded
    }
    if real := r.Header.Get("X-Real-IP"); real != "" {
        return real
    }
    host, _, err := net.SplitHostPort(r.RemoteAddr)
    if err != nil {
        return r.RemoteAddr
    }
    return host
}

4.11 限流与熔断的协同

// Package gateway 提供 API 网关限流与熔断协同
package gateway

import (
    "context"
    "net/http"
    "time"

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

// Gateway API 网关
type Gateway struct {
    // 入口限流(保护自身)
    globalLimiter *rate.Limiter
    userLimiters  *LRUIPRateLimiter

    // 下游熔断(保护下游)
    downstreamBreakers map[string]*CircuitBreaker
}

// HandleRequest 处理请求
func (g *Gateway) HandleRequest(downstream string, handler http.HandlerFunc) http.HandlerFunc {
    return func(w http.ResponseWriter, r *http.Request) {
        // 1. 入口限流
        if !g.globalLimiter.Allow() {
            respondRateLimited(w, 1, "系统繁忙")
            return
        }
        ip := extractIP(r)
        if !g.userLimiters.GetLimiter(ip).Allow() {
            respondRateLimited(w, 1, "请求过于频繁")
            return
        }

        // 2. 下游熔断
        breaker, ok := g.downstreamBreakers[downstream]
        if !ok {
            http.Error(w, "unknown downstream", http.StatusInternalServerError)
            return
        }
        if breaker.State() == gobreaker.StateOpen {
            http.Error(w, "downstream unavailable", http.StatusServiceUnavailable)
            return
        }

        // 3. 执行请求(被熔断器包装)
        _, err := breaker.Execute(r.Context(), func(ctx context.Context) (interface{}, error) {
            handler(w, r)
            return nil, nil
        })
        if err != nil {
            // 熔断或下游错误
            return
        }
    }
}

5. 算法对比与选型

5.1 限流算法对比

算法原理优点缺点适用场景
令牌桶固定速率补充令牌,请求消耗令牌允许突发、实现简单无法保证匀速API 网关、通用场景
漏桶固定速率处理请求,超出排队强制匀速、平滑无法应对突发消息队列、流处理
滑动窗口统计滑动时间窗口内请求数精确、无边界问题内存开销大精确限流、审计场景
固定窗口每个固定时间窗口内限制请求数实现简单、内存开销小边界问题(窗口切换瞬间可通过 2L)粗略限流
滑动日志记录每个请求时间戳最精确内存开销大、查询慢精确审计

5.2 限流粒度对比

粒度实现方式适用场景
全局限流单个 rate.Limiter 实例保护系统总容量
按 IP 限流sync.Map 维护每 IP 限流器防止单 IP 滥用
按用户限流以用户 ID 为 keyVIP 分级限流
按接口限流按路由 + 方法组合保护慢接口
按租户限流SaaS 多租户隔离租户配额管理
按地域限流按 GeoIP 限流异地容灾

5.3 熔断器实现对比

库优点缺点适用场景
sony/gobreaker简单、无依赖、API 清晰功能较少通用熔断
afex/hystrix-go功能完整、支持隔离已停维护、较重遗留系统
cep21/circuit高性能、可配置API 较复杂高性能场景
failsafe-go/failsafe支持重试、熔断、限流组合学习曲线综合容错

5.4 单机 vs 分布式限流

维度单机限流分布式限流
实现rate.LimiterRedis + Lua
性能O(1)O(1),亚微秒O(1)O(1),毫秒级(网络)
一致性单机内存Redis 强一致
部署单实例多实例共享
适用单机或近似限流精确全局限流
复杂性低高(需处理 Redis 故障)

6. 常见陷阱与反模式

6.1 陷阱一:按 IP 限流内存泄漏

问题:直接用 sync.Map 维护每 IP 限流器,IP 数量无限增长导致内存耗尽。

// 反例:无淘汰的按 IP 限流
limiters := sync.Map{}
func getLimiter(ip string) *rate.Limiter {
    if l, ok := limiters.Load(ip); ok {
        return l.(*rate.Limiter)
    }
    l := rate.NewLimiter(rate.Limit(10), 5)
    limiters.Store(ip, l) // 永不淘汰,内存泄漏!
    return l
}

修复:使用 LRU 淘汰或定期清理不活跃的限流器(见 5.4)。

6.2 陷阱二:限流粒度过粗

问题:仅全局限流,单个高频用户挤占其他用户配额。

// 反例:仅全局限流
globalLimiter := rate.NewLimiter(1000, 50)
// 一个恶意 IP 可以打满全局配额,导致正常用户被限流

修复:全局 + 用户双层限流。

6.3 陷阱三:分布式限流未考虑 Redis 故障

问题:Redis 不可用时限流完全失效或完全拒绝。

// 反例:Redis 故障时直接拒绝所有请求
allowed, err := limiter.Allow(ctx)
if err != nil {
    return false // Redis 故障 → 全部拒绝 → 服务不可用
}

修复:Redis 故障时降级为单机限流(允许一定误差)或快速失败(fail-open,避免完全不可用)。

allowed, err := limiter.Allow(ctx)
if err != nil {
    // Redis 故障:降级为本地限流
    return localLimiter.Allow(), nil
}
return allowed, nil

6.4 陷阱四:熔断器未区分业务错误与系统错误

问题:将 404、参数错误等业务错误计入熔断失败率,导致误熔断。

// 反例:所有错误都计入熔断
result, err := breaker.Execute(func() (interface{}, error) {
    resp, err := http.Get(url)
    if err != nil {
        return nil, err
    }
    if resp.StatusCode != 200 {
        return nil, fmt.Errorf("error: %d", resp.StatusCode) // 404 也算失败!
    }
    return nil, nil
})

修复:仅系统错误(5xx、超时、连接失败)触发熔断。

result, err := breaker.Execute(func() (interface{}, error) {
    resp, err := http.Get(url)
    if err != nil {
        return nil, err // 系统错误
    }
    if resp.StatusCode >= 500 {
        return nil, fmt.Errorf("server error: %d", resp.StatusCode) // 系统错误
    }
    if resp.StatusCode == 404 {
        return nil, ErrNotFound // 业务错误,不计入熔断
    }
    return nil, nil
})

6.5 陷阱五:限流响应缺少 Retry-After

问题:返回 429 但未告知客户端何时重试,导致客户端立即重试加剧拥塞。

// 反例:无 Retry-After
http.Error(w, "Too Many Requests", http.StatusTooManyRequests)

修复:返回 Retry-After 头。

w.Header().Set("Retry-After", "1")
http.Error(w, "Too Many Requests", http.StatusTooManyRequests)

6.6 陷阱六:熔断器恢复过快导致二次故障

问题:HalfOpen 状态下允许过多测试请求,下游未完全恢复又被打垮。

修复:MaxRequests 设为 1(逐个探测),Timeout 设为下游恢复时间的 2-3 倍。

6.7 陷阱七:Wait 阻塞导致 goroutine 泄漏

问题:在高并发下用 limiter.Wait(ctx) 阻塞等待,goroutine 堆积。

// 反例:每个请求阻塞等待
for i := 0; i < 10000; i++ {
    go func() {
        limiter.Wait(ctx) // 1 万个 goroutine 同时等待
        process()
    }()
}

修复:用 Allow 非阻塞模式 + 快速失败,或限制 goroutine 数量。

6.8 陷阱八:未对 WebSocket / 长连接限流

问题:WebSocket 连接建立后不再经过限流,长连接持续占用资源。

修复:对连接数限制(信号量),对消息速率限流。

6.9 陷阱九:滑动窗口边界问题

问题:固定窗口算法在窗口切换瞬间可通过 2L 请求(如 0.9s 时通过 L 个,1.0s 时新窗口又通过 L 个)。

修复:使用滑动窗口或令牌桶。

6.10 陷阱十:限流器未考虑时钟回拨

问题:NTP 时钟同步导致 time.Now() 回拨,令牌桶计算出现负的 elapsed。

修复:elapsed = max(0, now - last)。


7. 工程实践

7.1 限流组件设计

一个生产级限流组件应具备:

  1. 多维度限流:全局、用户、接口、租户。
  2. 多算法支持:令牌桶、滑动窗口、漏桶。
  3. 单机 + 分布式:单机优先,分布式精确。
  4. 降级策略:Redis 故障时降级单机。
  5. 监控指标:QPS、限流率、等待时间。
  6. 配置热更新:运行时调整限流参数。

7.2 配置化限流

# config.yaml 限流配置示例
rate_limit:
  global:
    rate: 10000        # 每秒 1 万请求
    burst: 1000
  user:
    rate: 100         # 每用户每秒 100 请求
    burst: 20
  api:
    - path: "/api/v1/search"
      method: "GET"
      rate: 50        # 搜索接口每秒 50
      burst: 10
    - path: "/api/v1/upload"
      method: "POST"
      rate: 5         # 上传接口每秒 5
      burst: 2
  redis:
    addr: "redis:6379"
    enabled: true     # 启用分布式限流
    fallback: true    # Redis 故障时降级单机

7.3 监控指标

// Package metrics 提供限流监控指标
package metrics

import (
    "github.com/prometheus/client_golang/prometheus"
    "github.com/prometheus/client_golang/prometheus/promauto"
)

var (
    // 请求总数
    requestTotal = promauto.NewCounterVec(prometheus.CounterOpts{
        Name: "rate_limit_requests_total",
        Help: "Total rate limit requests",
    }, []string{"path", "method", "result"}) // result: allowed, denied

    // 限流延迟
    limitDelay = promauto.NewHistogramVec(prometheus.HistogramOpts{
        Name:    "rate_limit_delay_seconds",
        Help:    "Rate limit delay",
        Buckets: prometheus.ExponentialBuckets(0.001, 2, 10),
    }, []string{"path", "method"})

    // 熔断器状态
    breakerState = promauto.NewGaugeVec(prometheus.GaugeOpts{
        Name: "circuit_breaker_state",
        Help: "Circuit breaker state (0=closed, 1=open, 2=half-open)",
    }, []string{"name"})

    // 熔断器统计
    breakerRequests = promauto.NewCounterVec(prometheus.CounterOpts{
        Name: "circuit_breaker_requests_total",
        Help: "Circuit breaker requests",
    }, []string{"name", "result"})
)

// RecordRequest 记录限流请求
func RecordRequest(path, method, result string) {
    requestTotal.WithLabelValues(path, method, result).Inc()
}

// RecordDelay 记录限流延迟
func RecordDelay(path, method string, delay float64) {
    limitDelay.WithLabelValues(path, method).Observe(delay)
}

// RecordBreakerState 记录熔断器状态
func RecordBreakerState(name string, state int) {
    breakerState.WithLabelValues(name).Set(float64(state))
}

7.4 限流与降级协同

// Package gateway 提供限流 + 熔断 + 降级的统一处理
package gateway

import (
    "context"
    "net/http"
)

// FallbackHandler 降级处理器
type FallbackHandler func(ctx context.Context, req *http.Request) (interface{}, error)

// GatewayConfig 网关配置
type GatewayConfig struct {
    EnableRateLimit bool
    EnableCircuit   bool
    EnableFallback  bool
    Fallbacks       map[string]FallbackHandler // 按 downstream 名称
}

// Handle 处理请求,统一限流 + 熔断 + 降级
func (g *Gateway) Handle(downstream string, handler http.HandlerFunc) http.HandlerFunc {
    return func(w http.ResponseWriter, r *http.Request) {
        // 1. 限流检查
        if g.cfg.EnableRateLimit && !g.checkRateLimit(r) {
            respondRateLimited(w, 1, "rate limited")
            return
        }

        // 2. 熔断检查
        if g.cfg.EnableCircuit {
            breaker := g.breakers[downstream]
            if breaker != nil && breaker.State() == gobreaker.StateOpen {
                // 3. 降级处理
                if g.cfg.EnableFallback {
                    if fb, ok := g.cfg.Fallbacks[downstream]; ok {
                        result, err := fb(r.Context(), r)
                        if err == nil {
                            respondJSON(w, result)
                            return
                        }
                    }
                }
                http.Error(w, "service unavailable", http.StatusServiceUnavailable)
                return
            }
        }

        // 4. 执行请求
        handler(w, r)
    }
}

7.5 测试限流逻辑

// Package ratelimit_test 提供限流器测试
package ratelimit_test

import (
    "context"
    "sync"
    "testing"
    "time"

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

// TestTokenBucket 测试令牌桶限流
func TestTokenBucket(t *testing.T) {
    limiter := rate.NewLimiter(10, 5) // 每秒 10,突发 5

    // 突发 5 个应全部通过
    for i := 0; i < 5; i++ {
        if !limiter.Allow() {
            t.Fatalf("请求 %d 应通过", i)
        }
    }

    // 第 6 个应被拒绝
    if limiter.Allow() {
        t.Fatal("第 6 个请求应被限流")
    }

    // 等待 100ms 后应能再获取 1 个令牌
    time.Sleep(100 * time.Millisecond)
    if !limiter.Allow() {
        t.Fatal("等待后应能通过")
    }
}

// TestSlidingWindow 测试滑动窗口限流
func TestSlidingWindow(t *testing.T) {
    // 使用内存版的滑动窗口
    sw := ratelimit.NewLocalSlidingWindow(10, time.Second)

    // 10 个请求应通过
    for i := 0; i < 10; i++ {
        if !sw.Allow() {
            t.Fatalf("请求 %d 应通过", i)
        }
    }

    // 第 11 个应被拒绝
    if sw.Allow() {
        t.Fatal("第 11 个请求应被限流")
    }

    // 等待 1 秒后窗口滑动,应能再次通过
    time.Sleep(time.Second)
    if !sw.Allow() {
        t.Fatal("窗口滑动后应通过")
    }
}

// TestCircuitBreaker 测试熔断器
func TestCircuitBreaker(t *testing.T) {
    cfg := Config{
        Name:          "test",
        MaxRequests:   1,
        Interval:      time.Second,
        Timeout:       100 * time.Millisecond,
        FailThreshold: 3,
    }
    cb := NewCircuitBreaker(cfg)

    // 触发 3 次失败,应进入 Open 状态
    for i := 0; i < 3; i++ {
        _, err := cb.Execute(context.Background(), func(ctx context.Context) (interface{}, error) {
            return nil, errors.New("fail")
        })
        if err == nil {
            t.Fatal("应返回错误")
        }
    }

    // 第 4 次应被熔断
    _, err := cb.Execute(context.Background(), func(ctx context.Context) (interface{}, error) {
        return "ok", nil
    })
    if err != ErrCircuitOpen {
        t.Fatalf("应返回熔断错误,实际:%v", err)
    }

    // 等待 Timeout 后进入 HalfOpen
    time.Sleep(150 * time.Millisecond)

    // 半开状态下允许 1 个请求
    result, err := cb.Execute(context.Background(), func(ctx context.Context) (interface{}, error) {
        return "recovered", nil
    })
    if err != nil || result != "recovered" {
        t.Fatalf("应恢复,实际:%v, %v", result, err)
    }
}

// TestConcurrency 测试并发场景下的限流
func TestConcurrency(t *testing.T) {
    limiter := rate.NewLimiter(100, 10)
    var allowed, denied int64
    var wg sync.WaitGroup

    // 并发 1000 个请求
    for i := 0; i < 1000; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            if limiter.Allow() {
                atomic.AddInt64(&allowed, 1)
            } else {
                atomic.AddInt64(&denied, 1)
            }
        }()
    }
    wg.Wait()

    // 突发 10 个应通过,其余被限流
    if allowed != 10 {
        t.Errorf("应允许 10 个,实际 %d", allowed)
    }
}

7.6 压测与基准测试

// BenchmarkRateLimiter 基准测试 rate.Limiter 性能
func BenchmarkRateLimiter(b *testing.B) {
    limiter := rate.NewLimiter(1000000, 1000000) // 高速率避免限流
    b.ResetTimer()
    b.RunParallel(func(pb *testing.PB) {
        for pb.Next() {
            limiter.Allow()
        }
    })
}

// BenchmarkSlidingWindow 基准测试滑动窗口
func BenchmarkSlidingWindow(b *testing.B) {
    sw := NewLocalSlidingWindow(1000000, time.Second)
    b.ResetTimer()
    b.RunParallel(func(pb *testing.PB) {
        for pb.Next() {
            sw.Allow()
        }
    })
}

8. 案例研究

8.1 案例一:电商 API 网关限流

某电商平台 API 网关限流策略:

  • 全局:10 万 QPS(保护网关总容量)。
  • 按用户:普通用户 100 QPS,VIP 1000 QPS。
  • 按接口:搜索 50 QPS,下单 5 QPS,秒杀 1 QPS。
  • 熔断:下游服务 5xx 失败率 > 10% 时熔断 30s。
  • 降级:熔断时返回缓存或默认推荐。
// 电商网关限流配置
config := CompositeLimitConfig{
    GlobalRate:     100000,
    GlobalBurst:    10000,
    UserRate:       100,
    UserBurst:      20,
    MaxConcurrency: 5000,
    RetryAfter:     1,
}

// VIP 用户提升限流
func getUserRate(user *User) rate.Limit {
    if user.IsVIP {
        return 1000
    }
    return 100
}

8.2 案例二:秒杀系统限流

秒杀场景特点:瞬间极高并发、库存有限。

  • 多级限流:CDN → 网关 → 应用 → 数据库。
  • 库存预热:库存加载到 Redis,用 DECR 原子扣减。
  • 请求排队:用 Kafka 削峰,控制消费速率。
  • 熔断:下游服务故障时返回”活动太火爆”。
// 秒杀限流:每秒仅允许 100 个请求进入下单流程
seckillLimiter := rate.NewLimiter(100, 50)

func SeckillHandler(w http.ResponseWriter, r *http.Request) {
    if !seckillLimiter.Allow() {
        respondJSON(w, map[string]interface{}{
            "code":    429,
            "message": "活动太火爆,请稍后再试",
        })
        return
    }

    // 检查库存(Redis 原子操作)
    stock, err := rdb.Incr(ctx, "seckill:stock").Result()
    if err != nil || stock > maxStock {
        respondJSON(w, map[string]interface{}{
            "code":    410,
            "message": "已售罄",
        })
        return
    }

    // 发送到 Kafka 异步处理
    produceOrder(stock)
}

8.3 案例三:微服务熔断保护

某微服务架构中,订单服务调用用户、商品、库存、支付四个下游服务:

  • 每个下游服务独立熔断器。
  • 熔断时返回缓存或默认值。
  • 监控告警:熔断次数 > 5 次/分钟触发告警。
type OrderService struct {
    userBreaker    *CircuitBreaker
    productBreaker *CircuitBreaker
    stockBreaker   *CircuitBreaker
    payBreaker     *CircuitBreaker
}

func (s *OrderService) CreateOrder(ctx context.Context, req *OrderReq) (*Order, error) {
    // 并行调用下游
    var wg sync.WaitGroup
    var user *User
    var product *Product
    var stockErr, userErr, productErr error

    wg.Add(2)
    go func() {
        defer wg.Done()
        result, err := s.userBreaker.Execute(ctx, func(ctx context.Context) (interface{}, error) {
            return userClient.Get(ctx, req.UserID)
        })
        if err != nil {
            userErr = err
            return
        }
        user = result.(*User)
    }()
    go func() {
        defer wg.Done()
        result, err := s.productBreaker.Execute(ctx, func(ctx context.Context) (interface{}, error) {
            return productClient.Get(ctx, req.ProductID)
        })
        if err != nil {
            productErr = err
            return
        }
        product = result.(*Product)
    }()
    wg.Wait()

    // 降级:用户服务熔断时使用默认用户
    if userErr == ErrCircuitOpen {
        user = &User{ID: req.UserID, Name: "unknown"}
    }

    // 创建订单
    return s.createOrder(ctx, user, product, req)
}

8.4 案例四:多租户 SaaS 限流

SaaS 平台按租户配额限流:

// 租户配置
type TenantConfig struct {
    TenantID string
    Plan     string // free, pro, enterprise
    RateLimit int   // 每秒请求数
}

// 租户限流器
type TenantLimiter struct {
    limiters sync.Map // map[tenantID]*rate.Limiter
    configs  map[string]TenantConfig
}

func (l *TenantLimiter) Allow(tenantID string) bool {
    cfg, ok := l.configs[tenantID]
    if !ok {
        return false // 未知租户
    }

    limiter, _ := l.limiters.LoadOrStore(tenantID,
        rate.NewLimiter(rate.Limit(cfg.RateLimit), cfg.RateLimit))
    return limiter.(*rate.Limiter).Allow()
}

8.5 案例五:Netflix 自适应限流思想

Netflix Concurrency Limits 的核心思想:

  1. 测量延迟:记录每次请求的延迟。
  2. 计算 Limit:基于 Little’s Law:L=throughput×tolerantLatency1L = \frac{\text{throughput} \times \text{tolerantLatency}}{1}
  3. AIMD 调整:
    • 增加(Additive Increase):低延迟时缓慢增加。
    • 减少(Multiplicative Decrease):高延迟时快速减少。
// 简化版 Netflix 自适应限流
type NetflixAdaptiveLimiter struct {
    mu              sync.Mutex
    currentLimit    int
    minLimit        int
    maxLimit        int
    latencyThreshold time.Duration
    queue           *Queue // 滑动窗口队列记录延迟
}

func (l *NetflixAdaptiveLimiter) OnSuccess(latency time.Duration) {
    l.mu.Lock()
    defer l.mu.Unlock()

    if latency < l.latencyThreshold {
        // 低延迟:加性增加
        if l.currentLimit < l.maxLimit {
            l.currentLimit++
        }
    }
}

func (l *NetflixAdaptiveLimiter) OnFailure(latency time.Duration) {
    l.mu.Lock()
    defer l.mu.Unlock()

    if latency > l.latencyThreshold {
        // 高延迟:乘性减少
        newLimit := l.currentLimit / 2
        if newLimit < l.minLimit {
            newLimit = l.minLimit
        }
        l.currentLimit = newLimit
    }
}

11. 扩展阅读

11.1 算法理论

  • Little’s Law:排队论基础,L=λWL = \lambda W,用于计算自适应限流。
  • AIMD(Additive Increase, Multiplicative Decrease):TCP 拥塞控制的核心算法,被限流借鉴。
  • Network Calculus:网络演算,用于理论分析限流算法的性能边界。

11.2 工程实践

  • Sentinel(阿里巴巴):支持流控、熔断、热点限流、系统自适应限流。
  • Envoy:Service Mesh 代理,内置限流与熔断能力。
  • Kong:API 网关,插件化限流。
  • Istio:基于 Envoy 的限流策略。

11.3 自适应限流

  • Netflix Concurrency Limits:基于梯度2算法的自适应限流。
  • Google BBR:基于带宽与延迟的拥塞控制,思想可借鉴。
  • Twitter RPC Adaptive Concurrency:基于请求延迟的自适应。

11.4 监控与可观测性

  • Prometheus:监控限流指标(QPS、限流率、延迟)。
  • Grafana:可视化限流与熔断状态。
  • OpenTelemetry:分布式追踪,关联限流与请求链路。

11.5 Go 生态

  • golang.org/x/sync/semaphore:基于 channel 的信号量实现。
  • github.com/sourcegraph/conc:提供 pool、waitgroup 等并发原语。
  • github.com/uber-go/ratelimit:Uber 开源的漏桶限流器。

12. 附录

12.1 限流算法速查表

算法Go 实现复杂度突发匀速适用
令牌桶rate.LimiterO(1)O(1)允许否通用
漏桶uber-go/ratelimitO(1)O(1)否是流处理
滑动窗口自实现 / Redis ZSETO(k)O(k) / O(log⁡N)O(\log N)否否精确限流
固定窗口自实现O(1)O(1)部分否粗略
并发限制semaphoreO(1)O(1)--连接池

12.2 rate.Limiter API 速查

方法说明阻塞
Allow()立即判断否
AllowN(now, n)判断 n 个令牌否
Wait(ctx)阻塞等待是
WaitN(ctx, n)阻塞等待 n 个是
Reserve()预留令牌否(返回延迟)
ReserveN(now, n)预留 n 个否
SetLimit(r)动态调整速率-
SetBurst(b)动态调整容量-
Limit()获取当前速率-
Burst()获取当前容量-
Tokens()获取当前令牌数-

12.3 熔断器状态转换图

stateDiagram-v2
    [*] --> Closed
    Closed --> Open: 失败率超阈值
    Open --> HalfOpen: 等待 Timeout
    HalfOpen --> HalfOpen: 失败
    HalfOpen --> Closed: 成功 N 次
    HalfOpen --> Open: 成功

12.4 HTTP 限流响应头

Header说明示例
Retry-After重试等待秒数Retry-After: 1
X-RateLimit-Limit总配额X-RateLimit-Limit: 1000
X-RateLimit-Remaining剩余配额X-RateLimit-Remaining: 950
X-RateLimit-Reset配额重置时间戳X-RateLimit-Reset: 1625000000

12.5 限流配置参考值

场景全局 QPS单用户 QPS突发
API 网关10000-100000100-100010-100
Web 应用1000-1000010-1005-20
秒杀100-10001-101-5
搜索1000-1000010-1005-20
数据库100-100010-505-10

12.6 常见问题 FAQ

Q1:令牌桶与漏桶如何选择?

A:大多数场景选令牌桶(允许突发,更友好);强制匀速选漏桶(如消息消费、数据库写入)。

Q2:分布式限流必须用 Redis 吗?

A:不必须。也可用 etcd、ZooKeeper、Consul 的分布式锁。但 Redis 因其高性能与 Lua 脚本支持,是最常用方案。

Q3:熔断器的 Timeout 设多少合适?

A:设为下游服务平均恢复时间的 2-3 倍。如下游重启需 10s,Timeout 设为 20-30s。

Q4:如何实现优先级限流?

A:维护多个限流器,按优先级顺序检查;高优先级先消耗令牌,低优先级排队。

Q5:限流与熔断有何区别?

A:限流保护自身不被打垮(控制请求入口);熔断保护自身不被下游拖垮(控制调用出口)。两者互补。

12.7 调试技巧

  1. 查看限流决策日志:记录每次 Allow 的结果与原因。
  2. Prometheus 指标:监控 QPS、限流率、延迟分布。
  3. 压测验证:用 wrk 或 vegeta 压测,验证限流是否生效。
  4. 熔断器状态查询:暴露 /debug/breaker 接口查询熔断状态。

12.8 版本兼容性

Go 版本rate.Limitergobreaker备注
Go 1.13+支持支持基础功能
Go 1.18+支持泛型-泛型减少部分需求
Go 1.20+支持v2 支持泛型推荐使用
Go 1.22+支持v2当前推荐

goroutine 创建

基本写法:启动 goroutine go <函数调用>

// 使用 go 关键字启动并发协程
go doWork()

基本写法:匿名函数 goroutine go func() { ... }()

// 立即执行的匿名 goroutine
go func(msg string) {
    fmt.Println(msg)
}("hello")

基本写法:WaitGroup 等待协程 var wg sync.WaitGroup

// 使用 WaitGroup 等待一组协程完成
var wg sync.WaitGroup
for i := 0; i < 3; i++ {
    wg.Add(1)
    go func(n int) {
        defer wg.Done()
        fmt.Println(n)
    }(i)
}
wg.Wait()

基本写法:Go 1.22+ 循环变量安全 for i := range <n> { go func() { ... }() }

// Go 1.22+ 每次迭代变量地址独立,无需手动拷贝
var wg sync.WaitGroup
for i := range 3 {
    wg.Add(1)
    go func() {
        defer wg.Done()
        fmt.Println(i)
    }()
}
wg.Wait()

channel 创建

基本写法:无缓冲通道 make(chan <类型>)

// 无缓冲通道:发送和接收同步
ch := make(chan int)

基本写法:带缓冲通道 make(chan <类型>, <容量>)

// 带缓冲通道:缓冲区满前发送不阻塞
ch := make(chan string, 3)

基本写法:单向发送通道 chan<- <类型>

// 只能发送的通道类型
var sendOnly chan<- int = make(chan int)

基本写法:单向接收通道 <-chan <类型>

// 只能接收的通道类型
var recvOnly <-chan int = make(chan int)

基本写法:关闭通道 close(<通道>)

// 关闭通道,通知接收方不再有数据
ch := make(chan int, 2)
ch <- 1
ch <- 2
close(ch)

channel 收发

基本写法:发送数据 <通道> <- <值>

// 向通道发送数据
ch := make(chan int, 1)
ch <- 42

基本写法:接收数据 <-<通道>

// 从通道接收数据
v := <-ch

基本写法:接收数据并判断是否关闭 v, ok := <-<通道>

// ok 为 false 表示通道已关闭且无数据
v, ok := <-ch
if !ok {
    fmt.Println("通道已关闭")
}

基本写法:range 遍历通道 for v := range <通道> { ... }

// 持续接收直到通道关闭
for v := range ch {
    fmt.Println(v)
}

基本写法:select 多路复用 select { case ... : }

// select 随机选择一个就绪的 case 执行
select {
case v := <-ch1:
    fmt.Println("ch1:", v)
case v := <-ch2:
    fmt.Println("ch2:", v)
}

基本写法:select 超时控制 select { case ... : case <-time.After(<时长>): }

// 使用 time.After 实现超时
select {
case v := <-ch:
    fmt.Println(v)
case <-time.After(2 * time.Second):
    fmt.Println("超时")
}

基本写法:select 非阻塞 select { case ...: default: }

// default 使 select 非阻塞
select {
case v := <-ch:
    fmt.Println(v)
default:
    fmt.Println("无数据")
}

通道容量与长度

基本写法:获取通道缓冲区容量 cap(<通道>)

// 返回通道缓冲区大小
ch := make(chan int, 5)
fmt.Println(cap(ch)) // 5

基本写法:获取通道中元素个数 len(<通道>)

// 返回通道中当前排队元素数
ch := make(chan int, 3)
ch <- 1
ch <- 2
fmt.Println(len(ch)) // 2

常用并发模式

换行写法:Worker Pool 工作池 func worker(id int, jobs <-chan int, results chan<- int)

// 固定数量 worker 消费任务队列
func worker(id int, jobs <-chan int, results chan<- int) {
    for j := range jobs {
        results <- j * 2
    }
}
jobs := make(chan int, 10)
results := make(chan int, 10)
for w := 1; w <= 3; w++ {
    go worker(w, jobs, results)
}
for j := 1; j <= 5; j++ {
    jobs <- j
}
close(jobs)

换行写法:扇入合并多通道 func fanIn(<通道1>, <通道2> <-chan <类型>) <-chan <类型>

// 将多个通道合并为一个
func fanIn(ch1, ch2 <-chan int) <-chan int {
    out := make(chan int)
    go func() { for v := range ch1 { out <- v } }()
    go func() { for v := range ch2 { out <- v } }()
    return out
}

换行写法:信号量限制并发数 sem := make(chan struct{}, <最大并发>)

// 使用带缓冲通道作为信号量控制并发
sem := make(chan struct{}, 3)
var wg sync.WaitGroup
for i := 0; i < 10; i++ {
    wg.Add(1)
    sem <- struct{}{}
    go func(n int) {
        defer wg.Done()
        defer func() { <-sem }()
        doWork(n)
    }(i)
}
wg.Wait()

换行写法:通知退出信号 quit := make(chan struct{})

// 使用空结构体通道作为退出信号
quit := make(chan struct{})
go func() {
    for {
        select {
        case <-quit:
            return
        default:
            doWork()
        }
    }
}()
close(quit)

Mutex 互斥锁

基本写法:互斥锁 var mu sync.Mutex

// 互斥锁保护共享资源
var mu sync.Mutex
count := 0
mu.Lock()
count++
mu.Unlock()

基本写法:读写锁 var rw sync.RWMutex

// 读写锁:允许多读单写
var rw sync.RWMutex
rw.RLock()
v := readData()
rw.RUnlock()

换行写法:使用 defer 解锁 mu.Lock(); defer mu.Unlock()

// defer 确保锁一定被释放
func update() {
    mu.Lock()
    defer mu.Unlock()
    doUpdate()
}

基本写法:tryLock 非阻塞加锁 mu.TryLock()

// 尝试加锁,失败返回 false
if mu.TryLock() {
    defer mu.Unlock()
    doWork()
} else {
    fmt.Println("加锁失败")
}

Once 单次执行

基本写法:sync.Once 确保只执行一次 var once sync.Once

// 单例模式初始化
var (
    instance *Config
    once     sync.Once
)
func GetConfig() *Config {
    once.Do(func() {
        instance = loadConfig()
    })
    return instance
}

atomic 原子操作

基本写法:原子加法 atomic.AddInt64(&<变量>, <值>)

// 原子整数加法
var count int64
atomic.AddInt64(&count, 1)

基本写法:原子读取 atomic.LoadInt64(&<变量>)

// 原子读取值
v := atomic.LoadInt64(&count)

基本写法:原子存储 atomic.StoreInt64(&<变量>, <值>)

// 原子存储值
atomic.StoreInt64(&count, 100)

基本写法:原子比较交换 atomic.CompareAndSwapInt64(&<变量>, <旧值>, <新值>)

// CAS 操作:值匹配旧值才更新
ok := atomic.CompareAndSwapInt64(&count, 10, 20)

基本写法:Go 1.19+ 原子类型 var <变量> atomic.Int64

// 使用类型安全的原子类型
var n atomic.Int64
n.Add(1)
n.Store(100)
fmt.Println(n.Load())

基本写法:原子指针 atomic.Pointer[<类型>]

// Go 1.19+ 泛型原子指针
var p atomic.Pointer[Config]
p.Store(&Config{Name: "default"})
cfg := p.Load()

Go 1.23+ Timer 通道变更

基本写法:Go 1.23+ Timer 可被 GC 回收 t := time.NewTimer(<时长>)

// Go 1.23+ 未 Stop 的 Timer 可被垃圾回收
t := time.NewTimer(time.Second)
// 即使不调用 t.Stop(),失去引用后也会被 GC 回收

基本写法:Go 1.23+ Timer 通道无缓冲 cap(t.C) == 0

// Go 1.23+ Timer 通道容量为 0,使用非阻塞 select 轮询
t := time.NewTimer(time.Second)
select {
case <-t.C:
    fmt.Println("超时")
default:
    fmt.Println("未超时")
}