Zer0e's Blog

AI其五:如何设计用量限流系统

字数统计: 9.6k阅读时长: 45 min
2026/07/25 Share

前言

最近在梳理某大模型平台的限流策略,因此也借此梳理下系统有哪些限流手段,最常见的就是用量限流了,例如我们买codex、各种coding plan都有一个几个小时的窗口限制,防止你一次性把token消耗完。
仔细想想,这个既简单又很复杂,也许是因为我没设计过这种系统,这次正好让AI帮我梳理下,但我发现他写的有点搓,我反反复复还对话了很多次,还手动修改了一部分。
注:此文代码均为AI生成,且未提供任何材料辅助AI生成。

正文

用量限流系统是 API 网关中的核心组件,用于控制用户在特定时间窗口内的请求量或 Token 消耗量。对于大语言模型服务而言,限流不仅是资源保护的手段,更是成本控制公平分配服务质量保障的基础设施。

本文将从需求分析、架构设计、技术选型到代码实现,完整讲解如何设计一个支持近 5 小时、近一周、近一月三个时间窗口的网关用量限流系统。

需求分析:为什么需要限流系统?

大模型服务的资源瓶颈

大语言模型的推理成本远高于传统 API:

1
2
3
4
5
6
7
8
9
传统 REST API:
- 单次请求成本: ~$0.0001
- 响应时间: 10-50ms
- 资源消耗: CPU 密集型

大模型 API:
- 单次请求成本: ~$0.01-0.10(取决于 Token 数)
- 响应时间: 500ms-5s
- 资源消耗: GPU 密集型(显存 + 计算)

核心问题:

flowchart TD
    A[大模型服务特点] --> B1["GPU 算力昂贵"]
    A --> B2["推理延迟高"]
    A --> B3["Token 成本按量计费"]
    
    B1 --> C["资源有限"]
    B2 --> C
    B3 --> C
    
    C --> D1["需要控制并发"]
    C --> D2["需要限制用量"]
    C --> D3["需要公平分配"]
    
    style A fill:#ffe0b2
    style C fill:#ff8a65,color:#fff
    style D1 fill:#66bb6a,color:#fff
    style D2 fill:#42a5f5,color:#fff
    style D3 fill:#ab47bc,color:#fff

限流的核心价值

维度 问题 限流的作用
资源保护 GPU 过载导致服务崩溃 防止突发流量压垮服务
成本控制 恶意调用或异常循环 避免巨额账单
公平分配 单一用户占用全部资源 多租户场景下的配额管理
服务质量 无限制导致延迟飙升 保障 SLA 和用户体验
计费基础 按用量定价 支撑 Tier 定价模型

限流 vs 熔断 vs 降级

flowchart TD
    A[服务保护机制] --> B1["限流 Rate Limiting"]
    A --> B2["熔断 Circuit Breaking"]
    A --> B3["降级 Fallback"]
    
    B1 --> C1["事前预防\n控制请求进入"]
    B2 --> C2["事中保护\n快速失败"]
    B3 --> C3["事后兜底\n返回降级响应"]
    
    style B1 fill:#66bb6a,color:#fff
    style B2 fill:#ffa726
    style B3 fill:#ef5350,color:#fff
  • 限流:控制流量进入的速度(事前)
  • 熔断:检测到异常后快速失败(事中)
  • 降级:返回简化响应保证可用性(事后)

限流系统核心概念

三个关键时间维度

大模型服务的用量限流通常采用多时间窗口策略:

flowchart TD
    A[用量限流] --> B1["近 5 小时\n短期突发控制"]
    A --> B2["近一周\n中期用量管理"]
    A --> B3["近一月\n长期配额规划"]
    
    B1 --> C1["粒度: 分钟级"]
    B1 --> D1["场景: 防刷、控并发"]
    
    B2 --> C2["粒度: 小时级"]
    B2 --> D2["场景: 周配额管理"]
    
    B3 --> C3["粒度: 天级"]
    B3 --> D3["场景: 月套餐限制"]
    
    style B1 fill:#42a5f5,color:#fff
    style B2 fill:#66bb6a,color:#fff
    style B3 fill:#ffa726
    style D1 fill:#e3f2fd
    style D2 fill:#e8f5e9
    style D3 fill:#fff3e0
时间窗口 控制粒度 典型场景 存储策略
近 5 小时 分钟级 短期突发控制、实时用量限制 Redis Sorted Set + Hash(滑动窗口)
近一周 小时级 周配额管理、按周计费 Redis(配额计数器)
近一月 天级 月套餐限制、月度预算控制 Redis + MySQL(配额计数器)

配额限流的核心概念

配额限流(Quota-based Rate Limiting)的核心是控制特定时间窗口内的总用量,而非单纯限制请求频率。

flowchart TD
    A[配额限流] --> B1["用量计量\nToken/请求数"]
    A --> B2["时间窗口\n5h/7d/30d"]
    A --> B3["配额上限\n动态配置"]
    
    B1 --> C["实时统计 + 超限拒绝"]
    B2 --> C
    B3 --> C
    
    C --> D1["短期:防突发"]
    C --> D2["中期:控预算"]
    C --> D3["长期:管套餐"]
    
    style A fill:#ff7043,color:#fff
    style C fill:#42a5f5,color:#fff
    style D1 fill:#66bb6a,color:#fff
    style D2 fill:#ffa726
    style D3 fill:#ab47bc,color:#fff

配额限流 vs 频率限流:

维度 频率限流(Rate Limiting) 配额限流(Quota Limiting)
控制对象 请求速率(QPS/RPM) 时间窗口内的总用量
典型场景 防刷、控并发 套餐管理、成本控制
时间维度 单一窗口(如 1 分钟) 多窗口组合(5h/7d/30d)
计量单位 请求次数 Token 数/请求数/Credits
重置策略 固定周期重置 按套餐周期/手动充值

限流算法选型

1. 固定窗口计数器(Fixed Window Counter)

最简单的限流算法:将时间划分为固定窗口,统计每个窗口内的请求数。

flowchart TD
    A[请求到达] --> B[确定当前窗口]
    B --> C[窗口内计数 +1]
    C --> D{计数 > 阈值?}
    D -->|是| E["❌ 拒绝请求"]
    D -->|否| F["✅ 放行请求"]
    
    style A fill:#bbdefb
    style D fill:#ffe0b2
    style E fill:#ef5350,color:#fff
    style F fill:#66bb6a,color:#fff

实现示例(Go):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
type FixedWindowLimiter struct {
mu sync.Mutex
window time.Duration
limit int
count int
windowStart time.Time
}

func (l *FixedWindowLimiter) Allow() bool {
l.mu.Lock()
defer l.mu.Unlock()

now := time.Now()
if now.Sub(l.windowStart) >= l.window {
// 窗口过期,重置
l.windowStart = now
l.count = 0
}

if l.count >= l.limit {
return false
}

l.count++
return true
}

缺陷:边界突变问题

1
2
3
4
5
时间轴:  |---窗口1---|---窗口2---|---窗口3---|
请求: | 100 | 100 | 100 |

问题: 在窗口交界处,可能出现 2 倍突发流量
窗口1末尾 100 次 + 窗口2开头 100 次 = 200 次/短时间

2. 滑动窗口计数器(Sliding Window Counter)

将窗口细分为多个子窗口,通过加权计算实现平滑过渡。

flowchart TD
    A[请求到达] --> B[计算当前时间戳]
    B --> C[移除过期子窗口]
    C --> D[聚合当前窗口内所有子窗口计数]
    D --> E{总数 > 阈值?}
    E -->|是| F["❌ 拒绝"]
    E -->|否| G["✅ 放行 + 当前子窗口 +1"]
    
    style A fill:#e3f2fd
    style D fill:#ffe0b2
    style F fill:#ef5350,color:#fff
    style G fill:#66bb6a,color:#fff

滑动窗口 vs 固定窗口对比:

flowchart LR
    subgraph Fixed["固定窗口"]
        A1["|===|===|===|"]
        A2["计数突变 ❌"]
    end
    
    subgraph Sliding["滑动窗口"]
        B1["|==|==|==|==|==|"]
        B2["平滑过渡 ✅"]
    end
    
    style Fixed fill:#ffebee
    style Sliding fill:#e8f5e9
    style A2 fill:#ffcdd2
    style B2 fill:#c8e6c9

3. 令牌桶算法(Token Bucket)

以固定速率生成令牌,请求需要消耗令牌才能通过,允许一定程度的突发。

flowchart TD
    A[令牌生成器] -->|"恒定速率 r"| B["令牌桶\n容量: N"]
    C[请求到达] --> D{桶中有令牌?}
    B --> D
    D -->|是| E["消耗 1 个令牌\n✅ 放行"]
    D -->|否| F["❌ 拒绝或排队"]
    
    style A fill:#ab47bc,color:#fff
    style B fill:#42a5f5,color:#fff
    style D fill:#ffe0b2
    style E fill:#66bb6a,color:#fff
    style F fill:#ef5350,color:#fff

核心公式:

当前令牌数计算:

$$\text{tokens} = \min(\text{capacity}, \text{old_tokens} + (\text{now} - \text{last_time}) \times \text{rate})$$

Go 实现:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
type TokenBucket struct {
mu sync.Mutex
rate float64 // 令牌生成速率(个/秒)
capacity float64 // 桶容量
tokens float64 // 当前令牌数
lastTime time.Time // 上次更新时间
}

func (tb *TokenBucket) Allow() bool {
tb.mu.Lock()
defer tb.mu.Unlock()

now := time.Now()
elapsed := now.Sub(tb.lastTime).Seconds()
tb.lastTime = now

// 补充令牌
tb.tokens = math.Min(tb.capacity, tb.tokens+elapsed*tb.rate)

if tb.tokens >= 1.0 {
tb.tokens--
return true
}
return false
}

4. 漏桶算法(Leaky Bucket)

与令牌桶相反,漏桶以固定速率流出请求,平滑突发流量。

flowchart TD
    A[请求流入] --> B["水桶\n容量: N"]
    B -->|"恒定速率 r"| C[请求流出]
    B --> D{桶满?}
    D -->|是| E["❌ 溢出丢弃"]
    D -->|否| F["✅ 入桶等待"]
    
    style A fill:#bbdefb
    style B fill:#ffa726,color:#fff
    style C fill:#66bb6a,color:#fff
    style E fill:#ef5350,color:#fff

算法对比总结

算法 特点 适用场景 是否允许突发
固定窗口 实现简单,边界突变 简单场景
滑动窗口 精确统计,存储开销大 精确用量控制
令牌桶 允许突发,实现复杂 API 限流
漏桶 平滑输出,延迟增加 流量整形

多时间窗口架构设计

核心挑战

同时维护 5 小时、7 天、30 天三个窗口,需要解决:

  1. 存储效率:细粒度数据存储成本高
  2. 查询性能:多窗口聚合查询延迟
  3. 原子性:分布式环境下计数一致性
  4. 过期清理:自动清理过期数据
flowchart TD
    A[多窗口限流挑战] --> B1["存储效率\n细粒度 vs 成本"]
    A --> B2["查询性能\n聚合延迟"]
    A --> B3["原子性\n分布式一致"]
    A --> B4["过期清理\n自动化管理"]
    
    B1 --> C["架构设计"]
    B2 --> C
    B3 --> C
    B4 --> C
    
    style A fill:#ff7043,color:#fff
    style C fill:#42a5f5,color:#fff

数据结构设计

配额限流系统根据时间窗口特性,采用不同的数据结构:

5 小时窗口:Sorted Set + Hash(滑动窗口)

5 小时窗口需要动态滑动,任何时刻查询都是过去 5 小时的用量,因此使用时间分桶方案。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
Key: rate_limit:user:{user_id}:5h:index
Value: Sorted Set
- Member: "2026-07-25-14:30" (桶名,30分钟粒度)
- Score: 1721894400 (桶起始时间戳)

Key: rate_limit:user:{user_id}:5h:data
Value: Hash
- Field: "2026-07-25-14:30"
- Value: 15000 (该桶的总 Token 数)

查询近 5 小时用量:
1. ZRANGEBYSCORE index (now-5*3600) now
→ 获取时间范围内的所有桶名
2. HMGET data [桶1, 桶2, ...]
→ 批量获取每个桶的 Token 数
3. 聚合求和 = 总 Token 用量

写入新请求:
1. 计算当前桶名: "2026-07-25-14:30"
2. HINCRBY data "2026-07-25-14:30" {token_cost}
3. ZADD index {timestamp} "2026-07-25-14:30" (幂等操作)

自动过期:
ZREMRANGEBYSCORE index 0 (now-5*3600)
→ 删除过期桶的索引(配合定时任务清理 Hash 数据)

优势:

  • ✅ 真正的滑动窗口,任何时刻查询都是过去 5 小时
  • ✅ 自动过期机制(ZREMRANGEBYSCORE)
  • ✅ 按桶聚合,内存占用可控(5小时/30分钟 = 10个桶)
  • ✅ 支持用量趋势分析

劣势:

  • ⚠️ 需要两个数据结构配合(Sorted Set + Hash)
  • ⚠️ 查询时需要聚合计算(但 10 个桶的聚合 < 1ms)

7 天/30 天窗口:配额计数器

长周期窗口不需要精确滑动,使用简单计数器即可。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
Key: rate_limit:user:{user_id}:7d:quota
Value: 100000 (配额上限,Token 数)

Key: rate_limit:user:{user_id}:7d:used
Value: 35000 (已使用量,Token 数)

Key: rate_limit:user:{user_id}:30d:quota
Value: 300000 (配额上限,Token 数)

Key: rate_limit:user:{user_id}:30d:used
Value: 150000 (已使用量,Token 数)

查询 7 天窗口:
1. GET rate_limit:user:{user_id}:7d:quota → 100000
2. GET rate_limit:user:{user_id}:7d:used → 35000
3. 剩余 = 100000 - 35000 = 65000

写入新请求:
INCRBY rate_limit:user:{user_id}:7d:used {token_cost}

自动过期:
EXPIRE rate_limit:user:{user_id}:7d:quota 604800 (7 天)
EXPIRE rate_limit:user:{user_id}:7d:used 604800 (7 天)

优势:

  • ✅ 数据结构简单,只需两个 Key
  • ✅ 查询快速(两次 GET)
  • ✅ 更新快速(INCRBY)
  • ✅ 自动过期(EXPIRE)

劣势:

  • ⚠️ 固定窗口,过期时突然清零
  • ⚠️ 无法追溯历史用量

多窗口检查流程

flowchart TD
    A[请求到达] --> B["提取用户ID、Token数"]
    B --> C["检查 5 小时窗口"]
    C --> D{5h 超限?}
    D -->|是| E["❌ 返回 429"]
    D -->|否| F["检查 7 天窗口"]
    
    F --> G{7d 超限?}
    G -->|是| E
    G -->|否| H["检查 30 天窗口"]
    
    H --> I{30d 超限?}
    I -->|是| E
    I -->|否| J["✅ 放行请求"]
    
    J --> K["异步更新计数器"]
    
    style A fill:#bbdefb
    style E fill:#ef5350,color:#fff
    style J fill:#66bb6a,color:#fff
    style K fill:#ffe0b2

配额存储与同步架构设计

核心问题:配额数据如何存储?

配额限流系统面临一个关键架构挑战:

1
2
3
4
5
用户购买配额 → 数据存储在哪里?
- MySQL?持久化强,但查询慢
- Redis?查询快,但可能丢失

Redis 缓存过期后,如何与 MySQL 同步?

方案一:Redis 主存储 + MySQL 持久化备份(推荐)

核心思路:配额数据以 Redis 为准,MySQL 仅做持久化和审计。

flowchart TD
    A[用户购买配额] --> B[写入 MySQL]
    B --> C[同步写入 Redis]
    
    D[请求到达] --> E[查询 Redis 配额]
    E --> F{配额充足?}
    F -->|是| G[放行 + Redis 扣减]
    F -->|否| H[拒绝]
    
    I[定时任务] --> J[Redis 用量同步到 MySQL]
    
    style B fill:#42a5f5,color:#fff
    style E fill:#66bb6a,color:#fff
    style G fill:#66bb6a,color:#fff
    style J fill:#ffa726

实现细节:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
// 1. 购买配额时:双写 MySQL + Redis
func PurchaseQuota(userID string, quota int, period string) error {
// 写入 MySQL(持久化)
_, err := db.Exec(`
INSERT INTO user_quotas (user_id, quota, period, created_at)
VALUES (?, ?, ?, NOW())
ON DUPLICATE KEY UPDATE quota = quota + ?
`, userID, quota, period, quota)

if err != nil {
return err
}

// 同步写入 Redis(用于限流)
redisKey := fmt.Sprintf("quota:user:%s:%s", userID, period)
pipe := redisClient.Pipeline()
pipe.SetNX(ctx, redisKey, quota, getExpiry(period)) // 设置过期时间
pipe.Exec(ctx)

return nil
}

// 2. 请求限流时:只读 Redis(亚毫秒级)
func CheckQuota(userID, period string, cost int) bool {
redisKey := fmt.Sprintf("quota:user:%s:%s", userID, period)

// 使用 Lua 脚本原子操作
script := `
local current = redis.call("GET", KEYS[1]) or 0
current = tonumber(current)
if current >= tonumber(ARGV[1]) then
redis.call("DECRBY", KEYS[1], ARGV[1])
return 1
end
return 0
`

result, _ := redisClient.Eval(ctx, script, []string{redisKey}, cost).Int()
return result == 1
}

// 3. 定时同步(每小时执行)
// ⚠️ 注意:生产环境禁止使用 KEYS 命令!
func SyncQuotaToMySQL() {
// ✅ 使用 SCAN 替代 KEYS,避免阻塞 Redis
var cursor uint64
pattern := "rate_limit:user:*:*:used"

for {
// 每次只扫描 100 个 Key
keys, newCursor, err := redisClient.Scan(ctx, cursor, pattern, 100).Result()
if err != nil {
log.Error("SCAN 失败:", err)
break
}

// 批量处理这批 Key
for _, key := range keys {
used, _ := redisClient.Get(ctx, key).Int()
// 解析 key 获取 userID 和 period
// Key 格式: rate_limit:user:{user_id}:{period}:used
parts := strings.Split(key, ":")
if len(parts) >= 5 {
userID := parts[2]
period := parts[3]

// 更新 MySQL
db.Exec(`
UPDATE user_quotas
SET used = ?, updated_at = NOW()
WHERE user_id = ? AND period = ?
`, used, userID, period)
}
}

cursor = newCursor
if cursor == 0 {
break // 扫描完成
}

// 避免过快扫描,给 Redis 喘息时间
time.Sleep(10 * time.Millisecond)
}
}

优势:

  • ✅ 高性能:限流检查只读 Redis(亚毫秒级)
  • ✅ 高可用:Redis 故障时可从 MySQL 恢复
  • ✅ 一致性:购买时双写,后续以 Redis 为准

劣势:

  • ⚠️ Redis 数据丢失风险(需配置持久化)
  • ⚠️ 需要定时同步任务

Redis 持久化配置:

1
2
3
4
5
6
7
# redis.conf
save 900 1 // 15 分钟内至少 1 次修改则保存
save 300 10 // 5 分钟内至少 10 次修改则保存
save 60 10000 // 1 分钟内至少 10000 次修改则保存

appendonly yes // 开启 AOF
appendfsync everysec // 每秒同步一次

方案二:MySQL 主存储 + Redis 懒加载

核心思路:MySQL 是主存储,Redis 仅作为缓存,缓存未命中时从 MySQL 加载。

flowchart TD
    A[请求到达] --> B{Redis 有缓存?}
    B -->|是| C[使用 Redis 配额]
    B -->|否| D[查询 MySQL]
    
    D --> E[加载配额到 Redis]
    E --> C
    
    C --> F{配额充足?}
    F -->|是| G[放行 + Redis 扣减]
    F -->|否| H[拒绝]
    
    G --> I[异步批量同步到 MySQL]
    
    style B fill:#ffe0b2
    style D fill:#42a5f5,color:#fff
    style G fill:#66bb6a,color:#fff
    style I fill:#ffa726

实现细节:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
func CheckQuotaLazyLoad(userID, period string, cost int) bool {
redisKey := fmt.Sprintf("quota:user:%s:%s", userID, period)

// 1. 尝试从 Redis 获取
quota, err := redisClient.Get(ctx, redisKey).Int()

if err == redis.Nil {
// 缓存未命中,从 MySQL 加载
var userQuota UserQuota
db.Where("user_id = ? AND period = ?", userID, period).First(&userQuota)

// 计算剩余配额
remaining := userQuota.Quota - userQuota.Used
quota = remaining

// 写入 Redis(设置过期时间)
redisClient.Set(ctx, redisKey, quota, getExpiry(period))
}

// 2. 检查配额并扣减
if quota >= cost {
redisClient.DecrBy(ctx, redisKey, int64(cost))
return true
}

return false
}

// 异步同步(使用消息队列)
func SyncToMySQLAsync(userID, period string, used int) {
// 发送到消息队列
producer.Send(&QuotaSyncMessage{
UserID: userID,
Period: period,
Used: used,
Timestamp: time.Now(),
})
}

// 消费者批量更新
func ConsumeSyncMessages() {
for msg := range consumer {
db.Exec(`
INSERT INTO user_quotas (user_id, period, used, updated_at)
VALUES (?, ?, ?, NOW())
ON DUPLICATE KEY UPDATE used = GREATEST(used, VALUES(used)),
updated_at = NOW()
`, msg.UserID, msg.Period, msg.Used)
}
}

优势:

  • ✅ MySQL 是单一数据源,数据最可靠
  • ✅ 懒加载减少不必要的 Redis 写入
  • ✅ 异步批量同步,降低数据库压力

劣势:

  • ⚠️ 首次请求延迟较高(需要查询 MySQL)
  • ⚠️ Redis 和 MySQL 可能有短暂不一致

方案三:双写 + 定期校对(最强一致性)

核心思路:同时写入 Redis 和 MySQL,定期校对差异并修复。

flowchart TD
    A[购买配额] --> B[双写 MySQL + Redis]
    
    C[请求扣减] --> D[更新 Redis]
    D --> E[异步写入 MySQL]
    
    F[定时校对任务] --> G[对比 Redis vs MySQL]
    G --> H{数据一致?}
    H -->|是| I[记录日志]
    H -->|否| J[以 Redis 为准修复 MySQL]
    
    style B fill:#66bb6a,color:#fff
    style D fill:#42a5f5,color:#fff
    style J fill:#ef5350,color:#fff

校对任务实现:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
func ReconcileQuotas() {
// 1. 获取 Redis 中的所有配额
redisQuotas := make(map[string]int)
keys := redisClient.Keys(ctx, "quota:user:*:*").Val()

for _, key := range keys {
val, _ := redisClient.Get(ctx, key).Int()
redisQuotas[key] = val
}

// 2. 查询 MySQL 中的配额
var mysqlQuotas []UserQuota
db.Find(&mysqlQuotas)

// 3. 对比并修复
for _, quota := range mysqlQuotas {
redisKey := fmt.Sprintf("quota:user:%s:%s", quota.UserID, quota.Period)
redisUsed := redisQuotas[redisKey]

mysqlUsed := quota.Used
diff := redisUsed - mysqlUsed

if abs(diff) > threshold { // 差异超过阈值
log.Warnf("配额不一致: %s, Redis=%d, MySQL=%d, 差异=%d",
redisKey, redisUsed, mysqlUsed, diff)

// 以 Redis 为准修复 MySQL
db.Exec(`
UPDATE user_quotas
SET used = ?, updated_at = NOW(), reconciled_at = NOW()
WHERE user_id = ? AND period = ?
`, redisUsed, quota.UserID, quota.Period)
}
}
}

方案对比与选型建议

维度 方案一:Redis 主存储 方案二:懒加载 方案三:双写+校对
性能 ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐ ⭐⭐⭐⭐
一致性 ⭐⭐⭐ ⭐⭐ ⭐⭐⭐⭐⭐
复杂度 ⭐⭐ ⭐⭐⭐ ⭐⭐⭐⭐
数据安全性 ⭐⭐⭐ ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐⭐
适用场景 高并发、可容忍少量丢失 中等并发、强一致性要求 金融级、不能出错

推荐方案:

对于大模型配额限流系统,推荐 方案一(Redis 主存储)+ 持久化配置

1
2
3
// Cron 定时任务配置
cron.AddFunc("0 * * * *", SyncQuotaToMySQL) // 每小时同步一次
cron.AddFunc("0 0 * * *", GenerateQuotaReport) // 每天生成报表

这样既能保证高性能,又能确保数据不丢失!

⚠️ 生产环境注意事项:避免 KEYS 命令陷阱

在定时同步任务中,绝对不能使用 KEYS 命令

1
2
// ❌ 危险操作!生产环境禁止使用
keys := redisClient.Keys(ctx, "quota:user:*:*").Val()

为什么 KEYS 命令很危险?

问题 影响 严重程度
阻塞 Redis KEYS 是 O(N) 复杂度,扫描期间阻塞所有请求 🔴 严重
带宽打满 返回大量 Key 占用网络带宽 🔴 严重
内存峰值 一次性加载所有 Key 到内存 🟡 中等
超时崩溃 大数据量时可能超时,导致同步失败 🔴 严重

真实案例:

1
2
3
4
5
6
某公司 10 万用户,每个用户 3 个窗口 = 30 万 Key
使用 KEYS 命令同步时:
- Redis 阻塞 3-5 秒
- 限流服务超时,大量请求失败
- Redis 带宽打满(500MB/s)
- 触发 Redis 内存告警

✅ 正确方案 1:使用 SCAN 迭代器

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
func SyncQuotaToMySQL() {
var cursor uint64
pattern := "rate_limit:user:*:*:used"

for {
// 每次只扫描 100 个 Key(非阻塞)
keys, newCursor, err := redisClient.Scan(ctx, cursor, pattern, 100).Result()
if err != nil {
log.Error(err)
break
}

// 处理这批 Key
for _, key := range keys {
// 同步到 MySQL
}

cursor = newCursor
if cursor == 0 {
break // 扫描完成
}

// 给 Redis 喘息时间
time.Sleep(10 * time.Millisecond)
}
}

SCAN vs KEYS 对比:

维度 KEYS SCAN
阻塞 是(阻塞所有请求) 否(迭代器模式)
复杂度 O(N) O(1) 每次调用
带宽 一次性返回所有 分批返回
生产可用 ❌ 禁止使用 ✅ 推荐使用

✅ 正确方案 2:维护 Key 索引表(最佳实践)

在 MySQL 中维护配额索引,避免扫描 Redis:

1
2
3
4
5
6
7
8
9
CREATE TABLE quota_index (
id INT PRIMARY KEY AUTO_INCREMENT,
user_id VARCHAR(64),
period VARCHAR(10),
redis_key VARCHAR(128),
needs_sync BOOLEAN DEFAULT TRUE,
synced_at TIMESTAMP,
INDEX idx_needs_sync (needs_sync)
);
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
func SyncQuotaToMySQL() {
// 1. 从 MySQL 获取需要同步的用户
var quotas []QuotaIndex
db.Where("needs_sync = ?", true).
Limit(1000).Find(&quotas)

// 2. 精确读取 Redis,不需要扫描
for _, q := range quotas {
used, _ := redisClient.Get(ctx, q.RedisKey).Int()

// 更新 MySQL
db.Exec("UPDATE user_quotas SET used = ? WHERE ...", used)

// 标记已同步
db.Exec("UPDATE quota_index SET needs_sync = FALSE, synced_at = NOW() WHERE id = ?", q.ID)
}
}

✅ 正确方案 3:事件驱动架构(大规模推荐)

每次配额变化时发送事件,异步同步到 MySQL:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// 请求扣减时
func CheckQuota(...) {
// ... 检查通过,更新 Redis

// 发送同步事件到 MQ
producer.Send(&QuotaChangeEvent{
UserID: userID,
Period: period,
Used: newUsed,
})
}

// 消费者异步同步
func ConsumeQuotaChanges() {
for event := range consumer {
db.Exec("UPDATE user_quotas SET used = ? WHERE ...", event.Used)
}
}

生产环境建议:

用户规模 推荐方案 理由
< 1 万 SCAN 迭代器 实现简单,性能可接受
1 万 - 10 万 Key 索引表 精确控制,避免扫描
> 10 万 事件驱动 实时同步,完全解耦

网关层技术选型与集成

API 网关选型

flowchart TD
    A[API 网关选型] --> B1["Kong"]
    A --> B2["Nginx + Lua"]
    A --> B3["Envoy"]
    A --> B4["自研网关"]
    
    B1 --> C1["✅ 插件生态丰富\n✅ 内置限流插件\n❌ 性能中等"]
    B2 --> C2["✅ 性能极高\n✅ 灵活定制\n❌ 开发成本高"]
    B3 --> C3["✅ 云原生支持\n✅ 可观测性强\n❌ 配置复杂"]
    B4 --> C4["✅ 完全定制\n❌ 维护成本高"]
    
    style A fill:#ff7043,color:#fff
    style C1 fill:#e8f5e9
    style C2 fill:#e3f2fd
    style C3 fill:#fff3e0
    style C4 fill:#fce4ec

推荐架构:Kong + Redis + Lua

flowchart TD
    A[客户端请求] --> B[Kong API Gateway]
    B --> C["Lua 限流插件"]
    C --> D{限流检查}
    D -->|通过| E[转发到后端服务]
    D -->|拒绝| F["返回 429 Too Many Requests"]
    
    C --> G[Redis 集群]
    G --> H["存储用量计数"]
    
    E --> I[大模型推理服务]
    I --> J[返回响应]
    J --> B
    B --> K[更新 Redis 计数]
    
    style A fill:#bbdefb
    style B fill:#42a5f5,color:#fff
    style G fill:#66bb6a,color:#fff
    style I fill:#ffa726
    style F fill:#ef5350,color:#fff

请求处理流程

flowchart TD
    A[HTTP 请求到达] --> B["解析 Header\n提取 UserID, AppID"]
    B --> C["计算本次 Token 消耗\n(预估或历史均值)"]
    C --> D["查询 Redis 获取当前用量"]
    D --> E{多窗口检查}
    
    E -->|通过| F["放行请求"]
    E -->|拒绝| G["返回 429 + Retry-After"]
    
    F --> H["异步更新 Redis"]
    H --> I["推送监控指标"]
    
    G --> J["记录限流日志"]
    
    style A fill:#e3f2fd
    style E fill:#ffe0b2
    style F fill:#66bb6a,color:#fff
    style G fill:#ef5350,color:#fff
    style H fill:#ab47bc,color:#fff

Redis Lua 脚本(原子性保证)

使用 Lua 脚本保证”检查 + 更新”的原子性:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
-- 多窗口限流 Lua 脚本
local user_id = KEYS[1]
local token_cost = tonumber(ARGV[1])

-- 窗口配置
local windows = {
{key_suffix = "5h", limit = 10000, expire = 18000},
{key_suffix = "7d", limit = 100000, expire = 604800},
{key_suffix = "30d", limit = 300000, expire = 2592000}
}

for _, window in ipairs(windows) do
local key = "rate_limit:user:" .. user_id .. ":" .. window.key_suffix

-- 获取当前用量
local current = redis.call("GET", key) or 0
current = tonumber(current)

-- 检查是否超限
if current + token_cost > window.limit then
return {0, "limit_exceeded:" .. window.key_suffix}
end
end

-- 所有窗口检查通过,更新计数
for _, window in ipairs(windows) do
local key = "rate_limit:user:" .. user_id .. ":" .. window.key_suffix
redis.call("INCRBY", key, token_cost)
redis.call("EXPIRE", key, window.expire)
end

return {1, "allowed"}

Go 调用示例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
var rateLimitScript = redis.NewScript(`
-- 上面的 Lua 脚本内容
`)

func (l *RateLimiter) CheckAndConsume(ctx context.Context, userID string, tokenCost int) (bool, error) {
result, err := rateLimitScript.Run(ctx, l.redis,
[]string{userID},
tokenCost,
).Result()

if err != nil {
return false, err
}

res := result.([]interface{})
allowed := res[0].(int64) == 1
return allowed, nil
}

高并发场景优化方案

Redis 性能瓶颈

单节点 Redis 的 QPS 限制约为 10 万,高并发场景需要优化:

flowchart TD
    A[高并发挑战] --> B1["Redis 单点瓶颈"]
    A --> B2["网络延迟累积"]
    A --> B3["热 Key 问题"]
    
    B1 --> C["优化方案"]
    B2 --> C
    B3 --> C
    
    C --> D1["Pipeline 批量操作"]
    C --> D2["本地缓存 + 异步同步"]
    C --> D3["Redis 集群分片"]
    
    style A fill:#ff7043,color:#fff
    style C fill:#42a5f5,color:#fff
    style D1 fill:#66bb6a,color:#fff
    style D2 fill:#ffa726
    style D3 fill:#ab47bc,color:#fff

优化方案 1:Pipeline 批量操作

将多次 Redis 调用合并为一次网络往返:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
func (l *RateLimiter) CheckBatch(ctx context.Context, userIDs []string) map[string]bool {
pipe := l.redis.Pipeline()

cmds := make(map[string]*redis.IntCmd)
for _, userID := range userIDs {
key := fmt.Sprintf("rate_limit:user:%s:5h", userID)
cmds[userID] = pipe.Get(ctx, key)
}

_, err := pipe.Exec(ctx)
if err != nil && err != redis.Nil {
// 处理错误
}

results := make(map[string]bool)
for userID, cmd := range cmds {
count, _ := cmd.Int()
results[userID] = count < 10000
}

return results
}

优化方案 2:本地缓存 + Redis 二级架构

flowchart TD
    A[请求到达] --> B{本地缓存检查}
    B -->|命中| C["本地计数器判断"]
    B -->|未命中| D["查询 Redis"]
    
    C --> E{本地是否超限?}
    E -->|是| F["❌ 拒绝"]
    E -->|否| G["✅ 放行"]
    
    D --> H["写入本地缓存"]
    H --> G
    
    G --> I["本地计数器 +1"]
    I --> J["异步批量同步到 Redis"]
    
    style A fill:#e3f2fd
    style B fill:#ffe0b2
    style F fill:#ef5350,color:#fff
    style G fill:#66bb6a,color:#fff
    style J fill:#ab47bc,color:#fff

Go 实现:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
type LocalCacheLimiter struct {
mu sync.Mutex
localCount map[string]int // 本地计数
redisClient *redis.Client
syncInterval time.Duration
batchSize int
}

func (l *LocalCacheLimiter) Allow(userID string, tokenCost int) bool {
l.mu.Lock()
defer l.mu.Unlock()

current := l.localCount[userID]
if current + tokenCost > 10000 {
return false
}

l.localCount[userID] += tokenCost

// 异步同步到 Redis(定时任务)
go l.syncToRedis()

return true
}

func (l *LocalCacheLimiter) syncToRedis() {
// 批量同步逻辑
// 使用 Pipeline 提高性能
}

优化方案 3:Redis 集群分片

flowchart LR
    A[请求] --> B[网关层]
    B --> C1["Redis Shard 1\n用户 A, D, G"]
    B --> C2["Redis Shard 2\n用户 B, E, H"]
    B --> C3["Redis Shard 3\n用户 C, F, I"]
    
    style A fill:#e3f2fd
    style B fill:#42a5f5,color:#fff
    style C1 fill:#66bb6a,color:#fff
    style C2 fill:#66bb6a,color:#fff
    style C3 fill:#66bb6a,color:#fff

通过用户 ID Hash 分片,分散单点压力。

完整代码实现(Go)

基于 Redis + Go 的多窗口限流器

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
package ratelimit

import (
"context"
"fmt"
"time"

"github.com/go-redis/redis/v8"
"gorm.io/gorm"
)

// WindowConfig 窗口配置
type WindowConfig struct {
Name string // 窗口名称(5h, 7d, 30d)
Window time.Duration // 窗口时长
Expire time.Duration // Redis 过期时间
}

// UserQuota 用户配额模型(MySQL)
type UserQuota struct {
ID uint `gorm:"primaryKey"`
UserID string `gorm:"uniqueIndex:idx_user_period"`
Period string `gorm:"uniqueIndex:idx_user_period"` // 5h, 7d, 30d
Quota int // 配额上限(Token 数)
Used int // 已使用量
UpdatedAt time.Time
}

// MultiWindowRateLimiter 多窗口限流器
type MultiWindowRateLimiter struct {
redis *redis.Client
db *gorm.DB // MySQL 连接
windows []WindowConfig // 窗口配置(从数据库读取配额)
}

// NewMultiWindowRateLimiter 创建限流器
func NewMultiWindowRateLimiter(redis *redis.Client, db *gorm.DB) *MultiWindowRateLimiter {
return &MultiWindowRateLimiter{
redis: redis,
db: db,
windows: []WindowConfig{
{Name: "5h", Window: 5 * time.Hour, Expire: 6 * time.Hour},
{Name: "7d", Window: 7 * 24 * time.Hour, Expire: 8 * 24 * time.Hour},
{Name: "30d", Window: 30 * 24 * time.Hour, Expire: 31 * 24 * time.Hour},
},
}
}

// CheckResult 检查结果
type CheckResult struct {
Allowed bool
Remaining map[string]int // 各窗口剩余用量
Reason string
}

// Check 检查并消耗用量
func (l *MultiWindowRateLimiter) Check(
ctx context.Context,
userID string,
tokenCost int,
) (*CheckResult, error) {
remaining := make(map[string]int)

for _, w := range l.windows {
var used int
var err error

// 5 小时窗口使用滑动窗口查询
if w.Name == "5h" {
used, err = l.getSlidingWindowUsed(ctx, userID, "5h", 5*60)
} else {
// 7 天/30 天窗口使用简单计数器
quotaKey := fmt.Sprintf("rate_limit:user:%s:%s:quota", userID, w.Name)
usedKey := fmt.Sprintf("rate_limit:user:%s:%s:used", userID, w.Name)

// 从 Redis 读取配额上限(优先)
quota, quotaErr := l.redis.Get(ctx, quotaKey).Int()

if quotaErr == redis.Nil {
// Redis 缓存未命中,从 MySQL 加载
var userQuota UserQuota
dbErr := l.db.Where("user_id = ? AND period = ?", userID, w.Name).First(&userQuota).Error

if dbErr == gorm.ErrRecordNotFound {
return &CheckResult{
Allowed: false,
Reason: fmt.Sprintf("%s 窗口无配额", w.Name),
}, nil
} else if dbErr != nil {
return nil, fmt.Errorf("查询配额失败: %w", dbErr)
}

quota = userQuota.Quota
l.redis.Set(ctx, quotaKey, quota, w.Expire)
l.redis.Set(ctx, usedKey, userQuota.Used, w.Expire)
} else if quotaErr != nil {
return nil, fmt.Errorf("读取 Redis 失败: %w", quotaErr)
}

used, err = l.redis.Get(ctx, usedKey).Int()
}

if err != nil && err != redis.Nil {
return nil, fmt.Errorf("查询用量失败: %w", err)
}

// 从 MySQL 获取配额上限
var userQuota UserQuota
dbErr := l.db.Where("user_id = ? AND period = ?", userID, w.Name).First(&userQuota).Error

if dbErr == gorm.ErrRecordNotFound {
return &CheckResult{
Allowed: false,
Reason: fmt.Sprintf("%s 窗口无配额", w.Name),
}, nil
} else if dbErr != nil {
return nil, fmt.Errorf("查询配额失败: %w", dbErr)
}

// 检查是否超限
if used+tokenCost > userQuota.Quota {
return &CheckResult{
Allowed: false,
Reason: fmt.Sprintf("%s 窗口用量已耗尽(%d/%d)", w.Name, used, userQuota.Quota),
}, nil
}

remaining[w.Name] = userQuota.Quota - used
}

// 所有窗口检查通过,更新用量
for _, w := range l.windows {
if w.Name == "5h" {
// 5 小时窗口:更新时间分桶
l.recordSlidingWindowUsage(ctx, userID, "5h", tokenCost)
} else {
// 7 天/30 天窗口:简单计数器
usedKey := fmt.Sprintf("rate_limit:user:%s:%s:used", userID, w.Name)
l.redis.IncrBy(ctx, usedKey, int64(tokenCost))
}
}

return &CheckResult{
Allowed: true,
Remaining: remaining,
}, nil
}

// getSlidingWindowUsed 查询滑动窗口已使用量
func (l *MultiWindowRateLimiter) getSlidingWindowUsed(
ctx context.Context,
userID, window string,
minutes int,
) (int, error) {
indexKey := fmt.Sprintf("rate_limit:user:%s:%s:index", userID, window)
dataKey := fmt.Sprintf("rate_limit:user:%s:%s:data", userID, window)

now := time.Now()
start := now.Add(-time.Duration(minutes) * time.Minute)

// 1. 获取时间范围内的所有桶
buckets, err := l.redis.ZRangeByScore(ctx, indexKey, &redis.ZRangeBy{
Min: fmt.Sprintf("%d", start.Unix()),
Max: fmt.Sprintf("%d", now.Unix()),
}).Result()

if err != nil {
return 0, err
}

if len(buckets) == 0 {
return 0, nil
}

// 2. 批量获取每个桶的 Token 数
tokenValues, err := l.redis.HMGet(ctx, dataKey, buckets...).Result()
if err != nil {
return 0, err
}

// 3. 聚合求和
total := 0
for _, val := range tokenValues {
if val != nil {
count, _ := strconv.Atoi(val.(string))
total += count
}
}

return total, nil
}

// recordSlidingWindowUsage 记录滑动窗口用量
func (l *MultiWindowRateLimiter) recordSlidingWindowUsage(
ctx context.Context,
userID, window string,
tokenCost int,
) {
indexKey := fmt.Sprintf("rate_limit:user:%s:%s:index", userID, window)
dataKey := fmt.Sprintf("rate_limit:user:%s:%s:data", userID, window)

now := time.Now()

// 计算当前桶名(30分钟粒度)
bucketTime := now.Truncate(30 * time.Minute)
bucketName := bucketTime.Format("2006-01-02-15:04")

// 1. 累加 Token 数到 Hash
l.redis.HIncrBy(ctx, dataKey, bucketName, int64(tokenCost))

// 2. 更新 Sorted Set 索引(幂等)
l.redis.ZAdd(ctx, indexKey, &redis.Z{
Score: float64(bucketTime.Unix()),
Member: bucketName,
})

// 3. 清理过期桶
cutoff := now.Add(-time.Duration(5) * time.Hour)
l.redis.ZRemRangeByScore(ctx, indexKey, "0", fmt.Sprintf("%d", cutoff.Unix()))
}

// PurchaseQuota 购买配额(写入 MySQL + Redis)
func (l *MultiWindowRateLimiter) PurchaseQuota(
ctx context.Context,
userID string,
period string,
quota int,
) error {
// 1. 写入或更新 MySQL
var userQuota UserQuota
err := l.db.Where("user_id = ? AND period = ?", userID, period).First(&userQuota).Error

if err == gorm.ErrRecordNotFound {
// 新建配额
userQuota = UserQuota{
UserID: userID,
Period: period,
Quota: quota,
Used: 0,
}
l.db.Create(&userQuota)
} else if err != nil {
return fmt.Errorf("查询配额失败: %w", err)
} else {
// 更新配额
l.db.Model(&userQuota).Update("quota", userQuota.Quota+quota)
}

// 2. 重置 Redis 使用量
usedKey := fmt.Sprintf("rate_limit:user:%s:%s:used", userID, period)
l.redis.Set(ctx, usedKey, 0, getExpiryByPeriod(period))

return nil
}

// getExpiryByPeriod 根据窗口获取过期时间
func getExpiryByPeriod(period string) time.Duration {
switch period {
case "5h":
return 6 * time.Hour
case "7d":
return 8 * 24 * time.Hour
case "30d":
return 31 * 24 * time.Hour
default:
return 24 * time.Hour
}
}

使用示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
func main() {
// 1. 初始化连接
redisClient := redis.NewClient(&redis.Options{
Addr: "localhost:6379",
})

db, _ := gorm.Open(mysql.Open("dsn..."), &gorm.Config{})

// 2. 创建限流器
limiter := ratelimit.NewMultiWindowRateLimiter(redisClient, db)

// 3. 用户购买配额
ctx := context.Background()
limiter.PurchaseQuota(ctx, "user_123", "5h", 10000) // 5小时 1万 Token
limiter.PurchaseQuota(ctx, "user_123", "7d", 100000) // 7天 10万 Token
limiter.PurchaseQuota(ctx, "user_123", "30d", 300000) // 30天 30万 Token

// 4. 请求时检查配额
result, err := limiter.Check(ctx, "user_123", 1500) // 本次请求 1500 Token
if err != nil {
log.Fatal(err)
}

if !result.Allowed {
fmt.Printf("配额不足: %s\n", result.Reason)
// 返回 429 错误
} else {
fmt.Printf("配额充足,放行请求\n")
fmt.Printf("5h 剩余: %d\n", result.Remaining["5h"])
fmt.Printf("7d 剩余: %d\n", result.Remaining["7d"])
fmt.Printf("30d 剩余: %d\n", result.Remaining["30d"])
}
}

网关中间件集成

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
package middleware

import (
"net/http"
"strconv"

"your-project/ratelimit"
)

// RateLimitMiddleware 限流中间件
func RateLimitMiddleware(limiter *ratelimit.MultiWindowRateLimiter) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()

// 提取用户 ID
userID := r.Header.Get("X-User-ID")
if userID == "" {
http.Error(w, "Missing X-User-ID header", http.StatusUnauthorized)
return
}

// 预估 Token 消耗(或使用历史均值)
tokenCost := estimateTokenCost(r)

// 限流检查
result, err := limiter.Check(ctx, userID, tokenCost)
if err != nil {
http.Error(w, "Internal Server Error", http.StatusInternalServerError)
return
}

if !result.Allowed {
w.Header().Set("X-RateLimit-Reason", result.Reason)
w.Header().Set("Retry-After", "3600")
http.Error(w, "Rate Limit Exceeded", http.StatusTooManyRequests)
return
}

// 添加响应头
for window, remaining := range result.Remaining {
w.Header().Set(
fmt.Sprintf("X-RateLimit-Remaining-%s", window),
strconv.Itoa(remaining),
)
}

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

// 估算 Token 消耗
func estimateTokenCost(r *http.Request) int {
// 基于请求大小估算,或使用默认值
return 1000 // 示例:默认 1000 tokens
}

实际应用场景与配置

1. SaaS 平台的 Tier 定价

flowchart TD
    A[套餐类型] --> B1["免费版"]
    A --> B2["专业版"]
    A --> B3["企业版"]
    
    B1 --> C1["5h: 1000 tokens\n7d: 5000 tokens\n30d: 10000 tokens"]
    B2 --> C2["5h: 10000 tokens\n7d: 50000 tokens\n30d: 200000 tokens"]
    B3 --> C3["5h: 100000 tokens\n7d: 500000 tokens\n30d: 2000000 tokens"]
    
    style A fill:#ff7043,color:#fff
    style C1 fill:#ef5350,color:#fff
    style C2 fill:#ffa726
    style C3 fill:#66bb6a,color:#fff

配置示例:

1
2
3
4
5
6
7
8
9
10
11
12
type TierConfig struct {
Name string
Limit5h int
Limit7d int
Limit30d int
}

var Tiers = map[string]TierConfig{
"free": {Limit5h: 1000, Limit7d: 5000, Limit30d: 10000},
"pro": {Limit5h: 10000, Limit7d: 50000, Limit30d: 200000},
"enterprise": {Limit5h: 100000, Limit7d: 500000, Limit30d: 2000000},
}

2. API 市场的用量包管理

flowchart LR
    A[用户购买用量包] --> B["充值 100K tokens"]
    B --> C["余额: 100000"]
    C --> D["每次请求扣除"]
    D --> E{余额不足?}
    E -->|是| F["拒绝 + 提示充值"]
    E -->|否| G["放行 + 扣减余额"]
    
    style A fill:#e3f2fd
    style C fill:#66bb6a,color:#fff
    style F fill:#ef5350,color:#fff
    style G fill:#42a5f5,color:#fff

3. 企业内部的部门配额分配

1
2
3
4
5
企业总配额: 1000K tokens/月
├── 研发部: 400K tokens/月
├── 市场部: 300K tokens/月
├── 客服部: 200K tokens/月
└── 预留池: 100K tokens/月(动态分配)

4. 突发流量削峰

flowchart TD
    A[正常流量] --> B["1000 Tokens/min"]
    A --> C[突发流量]
    C --> D["10000 Tokens/min"]
    
    D --> E{配额检查}
    E -->|配额充足| F["2000 Tokens/min\n允许 2 倍突发"]
    E -->|配额不足| G["8000 Tokens/min\n返回 429"]
    
    style A fill:#e8f5e9
    style D fill:#ffebee
    style F fill:#66bb6a,color:#fff
    style G fill:#ef5350,color:#fff

通过令牌桶算法允许短期突发,同时保护长期配额。

监控告警体系设计

用量可视化

flowchart TD
    A[限流系统] --> B[数据采集]
    B --> C[时序数据库]
    C --> D[Dashboard]
    
    D --> E1["实时用量曲线"]
    D --> E2["各窗口剩余量"]
    D --> E3["限流触发次数"]
    D --> E4["用户 Top 排名"]
    
    style A fill:#42a5f5,color:#fff
    style C fill:#66bb6a,color:#fff
    style D fill:#ffa726

告警策略

告警级别 触发条件 动作
预警 用量达到 80% 发送邮件/消息通知
警告 用量达到 90% 发送短信 + 限流预警
严重 用量达到 100% 触发限流 + 通知运营
异常 短时间内大量限流 触发安全告警(可能是攻击)

Go 实现告警检查:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
func (m *Monitor) CheckAndAlert(userID string, usage int, limit int) {
ratio := float64(usage) / float64(limit)

if ratio >= 1.0 {
m.sendAlert(userID, "critical", "用量已耗尽")
} else if ratio >= 0.9 {
m.sendAlert(userID, "warning", "用量达到 90%")
} else if ratio >= 0.8 {
m.sendAlert(userID, "info", "用量达到 80%")
}
}

func (m *Monitor) sendAlert(userID, level, message string) {
// 集成告警系统(邮件、短信、Webhook)
alert := Alert{
UserID: userID,
Level: level,
Message: message,
Time: time.Now(),
}

m.alertService.Send(alert)
}

异常检测

flowchart TD
    A[限流日志] --> B[异常检测引擎]
    B --> C1["频次异常\n短时间大量请求"]
    B --> C2["模式异常\n非正常调用模式"]
    B --> C3["地域异常\n异常 IP 来源"]
    
    C1 --> D[触发安全策略]
    C2 --> D
    C3 --> D
    
    D --> E1["临时封禁 IP"]
    D --> E2["验证码验证"]
    D --> E3["人工审核"]
    
    style A fill:#e3f2fd
    style B fill:#ab47bc,color:#fff
    style D fill:#ef5350,color:#fff

系统局限性与改进方向

1. 配额精度的权衡

flowchart TD
    A[配额精度] --> B1["高精度\n按请求记录"]
    A --> B2["中精度\n配额计数器"]
    A --> B3["低精度\n固定周期重置"]
    
    B1 --> C1["✅ 可追溯\n❌ 存储开销大"]
    B2 --> C2["✅ 平衡\n✅ 性能高"]
    B3 --> C3["✅ 简单\n❌ 边界突变"]
    
    style A fill:#ff7043,color:#fff
    style C1 fill:#e3f2fd
    style C2 fill:#c8e6c9
    style C3 fill:#fff3e0

2. 跨地域分布式限流的挑战

  • 数据同步延迟:多地域 Redis 同步导致计数不一致
  • 网络分区:部分地域无法访问中心 Redis
  • 解决方案
    • 本地限流 + 全局聚合
    • CRDT(无冲突复制数据类型)
    • 异步最终一致性

3. 预估 Token 消耗的误差

1
2
3
4
5
6
7
8
9
10
问题: 请求时无法精确知道实际 Token 消耗

方案 1: 使用历史均值
→ 误差: ±20%

方案 2: 先放行,后扣减
→ 风险: 可能超限

方案 3: 保守估算 + 事后调整
→ 平衡: 安全但可能过度限流

4. 与计费系统的深度集成

限流系统需要与计费系统协同:

flowchart TD
    A[请求到达] --> B[限流系统]
    B --> C[计费系统]
    C --> D[实际 Token 消耗]
    D --> E[限流系统回写]
    E --> F[修正计数器]
    
    style A fill:#e3f2fd
    style B fill:#66bb6a,color:#fff
    style C fill:#42a5f5,color:#fff
    style E fill:#ab47bc,color:#fff

总结与最佳实践

大模型限流系统是 API 网关的核心基础设施,通过多维度、多时间窗口的组合策略,实现资源保护、成本控制和公平分配。

flowchart TD
    A[限流系统核心价值] --> B1["资源保护\n防止过载"]
    A --> B2["成本控制\n避免滥用"]
    A --> B3["公平分配\n多租户管理"]
    A --> B4["计费基础\n支撑商业化"]
    
    B1 --> C[高可用大模型服务]
    B2 --> C
    B3 --> C
    B4 --> C
    
    C --> D1["✅ 多窗口限流"]
    C --> D2["✅ Redis + Lua 原子操作"]
    C --> D3["✅ 高并发优化"]
    C --> D4["✅ 监控告警"]
    
    style A fill:#ff7043,color:#fff
    style C fill:#42a5f5,color:#fff
    style D1 fill:#66bb6a
    style D2 fill:#ab47bc
    style D3 fill:#26c6da
    style D4 fill:#ffa726

技术选型建议:

场景 推荐方案 原因
初创项目 Kong + Redis 插件 快速上线,生态成熟
高并发场景 自研网关 + Redis 集群 性能优先,灵活定制
云原生架构 Envoy + Redis 微服务友好,可观测性强
混合云 本地限流 + 全局聚合 降低跨地域延迟

架构要点回顾:

  • 多窗口组合:5h(分钟级)+ 7d(小时级)+ 30d(天级)
  • 原子操作:Redis Lua 脚本保证一致性
  • 高并发优化:Pipeline + 本地缓存 + 集群分片
  • 监控告警:实时可视化 + 多级预警

限流系统不仅是技术组件,更是业务策略的体现。合理设计限流规则,可以在保护资源的同时,为用户提供最佳服务体验。

参考资料

  1. Kong Documentation: Rate Limiting Plugin
  2. Redis Documentation: Rate Limiting Patterns
  3. Google SRE Book: Chapter 15 - Traffic Management
  4. Go-Redis Documentation: Pipeline and Lua Script
  5. Envoy Documentation: Local Rate Limit Filter
  6. Nginx Documentation: Limiting Request Rate
CATALOG
  1. 1. 前言
  2. 2. 正文
    1. 2.1. 需求分析:为什么需要限流系统?
      1. 2.1.1. 大模型服务的资源瓶颈
      2. 2.1.2. 限流的核心价值
      3. 2.1.3. 限流 vs 熔断 vs 降级
    2. 2.2. 限流系统核心概念
      1. 2.2.1. 三个关键时间维度
      2. 2.2.2. 配额限流的核心概念
    3. 2.3. 限流算法选型
      1. 2.3.1. 1. 固定窗口计数器(Fixed Window Counter)
      2. 2.3.2. 2. 滑动窗口计数器(Sliding Window Counter)
      3. 2.3.3. 3. 令牌桶算法(Token Bucket)
      4. 2.3.4. 4. 漏桶算法(Leaky Bucket)
      5. 2.3.5. 算法对比总结
    4. 2.4. 多时间窗口架构设计
      1. 2.4.1. 核心挑战
      2. 2.4.2. 数据结构设计
        1. 2.4.2.1. 5 小时窗口:Sorted Set + Hash(滑动窗口)
        2. 2.4.2.2. 7 天/30 天窗口:配额计数器
      3. 2.4.3. 多窗口检查流程
    5. 2.5. 配额存储与同步架构设计
      1. 2.5.1. 核心问题:配额数据如何存储?
      2. 2.5.2. 方案一:Redis 主存储 + MySQL 持久化备份(推荐)
      3. 2.5.3. 方案二:MySQL 主存储 + Redis 懒加载
      4. 2.5.4. 方案三:双写 + 定期校对(最强一致性)
      5. 2.5.5. 方案对比与选型建议
      6. 2.5.6. ⚠️ 生产环境注意事项:避免 KEYS 命令陷阱
    6. 2.6. 网关层技术选型与集成
      1. 2.6.1. API 网关选型
      2. 2.6.2. 推荐架构:Kong + Redis + Lua
      3. 2.6.3. 请求处理流程
      4. 2.6.4. Redis Lua 脚本(原子性保证)
    7. 2.7. 高并发场景优化方案
      1. 2.7.1. Redis 性能瓶颈
      2. 2.7.2. 优化方案 1:Pipeline 批量操作
      3. 2.7.3. 优化方案 2:本地缓存 + Redis 二级架构
      4. 2.7.4. 优化方案 3:Redis 集群分片
    8. 2.8. 完整代码实现(Go)
      1. 2.8.1. 基于 Redis + Go 的多窗口限流器
      2. 2.8.2. 使用示例
      3. 2.8.3. 网关中间件集成
    9. 2.9. 实际应用场景与配置
      1. 2.9.1. 1. SaaS 平台的 Tier 定价
      2. 2.9.2. 2. API 市场的用量包管理
      3. 2.9.3. 3. 企业内部的部门配额分配
      4. 2.9.4. 4. 突发流量削峰
    10. 2.10. 监控告警体系设计
      1. 2.10.1. 用量可视化
      2. 2.10.2. 告警策略
      3. 2.10.3. 异常检测
    11. 2.11. 系统局限性与改进方向
      1. 2.11.1. 1. 配额精度的权衡
      2. 2.11.2. 2. 跨地域分布式限流的挑战
      3. 2.11.3. 3. 预估 Token 消耗的误差
      4. 2.11.4. 4. 与计费系统的深度集成
    12. 2.12. 总结与最佳实践
    13. 2.13. 参考资料