vault backup: 2026-06-01 10:06:16
This commit is contained in:
+109
-72
@@ -9,38 +9,47 @@ create time: 2026-06-01 10:30
|
||||
|
||||
Redis 的每次命令执行都有网络往返(RTT)开销。对于需要执行多条命令的场景,**Pipeline** 可以把多个命令打包成一个批次发送,大幅降低 RTT 次数。这是限流性能优化的核心手段之一(F14 考点)。
|
||||
|
||||
> [!question]- 思考一下
|
||||
>
|
||||
> 假设你的应用和 Redis 在同一机房(RTT ≈ 0.5ms),每条命令平均处理时间 0.1ms。如果一次请求需要 10 条 Redis 命令,总延迟是多少?用 Pipeline 打包后呢?
|
||||
|
||||
## 一、为什么需要 Pipeline?
|
||||
|
||||
### 无 Pipeline:逐条发送
|
||||
|
||||
```
|
||||
应用 Redis
|
||||
│ ───INCR key1─────────────▶ │ (RTT #1)
|
||||
│ ◀──1────────────────────── │
|
||||
│ │
|
||||
│ ───EXPIRE key1 60─────────▶ │ (RTT #2)
|
||||
│ ◀──1────────────────────── │
|
||||
│ │
|
||||
│ ───INCR key2─────────────▶ │ (RTT #3)
|
||||
│ ◀──1────────────────────── │
|
||||
│ │
|
||||
│ ───EXPIRE key2 60─────────▶ │ (RTT #4)
|
||||
│ ◀──1────────────────────── │
|
||||
│ │
|
||||
总 RTT: 4 次,延迟 = 4 × network_latency
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant App as 应用
|
||||
participant R as Redis
|
||||
|
||||
App->>R: INCR key1
|
||||
Note over App,R: RTT #1
|
||||
R-->>App: 1
|
||||
App->>R: EXPIRE key1 60
|
||||
Note over App,R: RTT #2
|
||||
R-->>App: 1
|
||||
App->>R: INCR key2
|
||||
Note over App,R: RTT #3
|
||||
R-->>App: 1
|
||||
App->>R: EXPIRE key2 60
|
||||
Note over App,R: RTT #4
|
||||
R-->>App: 1
|
||||
|
||||
Note over App: 总 RTT: 4 次<br/>延迟 = 4 × network_latency
|
||||
```
|
||||
|
||||
### 有 Pipeline:批量发送
|
||||
|
||||
```
|
||||
应用 Redis
|
||||
│ ───INCR key1 │
|
||||
│ ───EXPIRE key1 60 │ (一次 TCP 发送)
|
||||
│ ───INCR key2 │
|
||||
│ ───EXPIRE key2 60─────────▶ │ (RTT #1)
|
||||
│ ◀──[1, 1, 1, 1]─────────── │
|
||||
│ │
|
||||
总 RTT: 1 次,延迟 = 1 × network_latency
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant App as 应用
|
||||
participant R as Redis
|
||||
|
||||
App->>+R: INCR key1<br/>EXPIRE key1 60<br/>INCR key2<br/>EXPIRE key2 60
|
||||
Note over App,R: 一次 TCP 发送 (Pipeline 打包)
|
||||
R-->>-App: [1, 1, 1, 1]
|
||||
|
||||
Note over App: 总 RTT: 1 次<br/>延迟 = 1 × network_latency
|
||||
```
|
||||
|
||||
**性能提升:** 如果网络 RTT 是 1ms,4 条命令从 4ms 降到 1ms——**节省了 75% 的网络延迟**。
|
||||
@@ -55,89 +64,110 @@ Redis 的每次命令执行都有网络往返(RTT)开销。对于需要执
|
||||
|------|----------|------------|
|
||||
| **原子性** | ❌ 每条命令独立执行 | ✅ EXEC 时整体执行 |
|
||||
| **取消支持** | ❌ 不能中途取消 | ✅ DISCARD 取消 |
|
||||
| **事务回滚** | ❌ 单条失败不影响其他 | ⚠️ EXEC 失败则全部不执行 |
|
||||
| **事务回滚** | ❌ 单条失败不影响其他 | ⚠️ 编译错误时不执行,运行时错误照常返回 |
|
||||
| **嵌套管道** | ❌ 不能在事务内用 Pipeline | ❌ 不支持嵌套 |
|
||||
| **Watch 支持** | ❌ 无 | ✅ 配合 WATCH 实现乐观锁 |
|
||||
|
||||
### 关键区别:原子性
|
||||
|
||||
> [!warning]- Redis 事务的"坑"
|
||||
>
|
||||
> Redis 的 MULTI/EXEC **不提供回滚语义**!如果某条命令在 EXEC 时因为类型错误等运行时异常失败,其他命令仍会执行。它只保证 EXEC 后的命令不被其他客户端打断——也就是**串行化**,而非传统数据库的 ACID 事务。
|
||||
|
||||
```go
|
||||
// Pipeline:每条命令独立执行,前一条成功后面失败也照常返回
|
||||
pipe := rdb.Pipeline()
|
||||
pipe.Incr(ctx, "key1") // 成功
|
||||
pipe.Expire(ctx, "key1", 60) // 即使这步出错,Incr 结果仍会返回
|
||||
pipe.Incr(ctx, "key1") // 成功写入
|
||||
pipe.Expire(ctx, "key1", 60) // 即使这步出错,Incr 结果仍会返回
|
||||
results, _ := pipe.Exec(ctx)
|
||||
|
||||
// MULTI/EXEC:EXEC 时所有命令作为一个整体执行
|
||||
pipe2 := rdb.TxPipeline() // TxPipeline = MULTI/EXEC 包装
|
||||
// TxPipeline:MULTI/EXEC 包装,命令串行执行(不被其他客户端插入)
|
||||
// 注意:不是真正的原子回滚!
|
||||
pipe2 := rdb.TxPipeline() // = MULTI ... EXEC 包装
|
||||
pipe2.Incr(ctx, "key1")
|
||||
pipe2.Expire(ctx, "key1", 60)
|
||||
results2, _ := pipe2.Exec(ctx)
|
||||
```
|
||||
|
||||
> [!tip]- Go-Redis 的 API 区分
|
||||
### 如何选择:Pipeline 还是 TxPipeline?
|
||||
|
||||
> [!tip]- Go-Redis 的 API 选择指南
|
||||
>
|
||||
> | API | 对应行为 |
|
||||
> |-----|---------|
|
||||
> | `rdb.Pipeline()` | 普通 Pipeline(非原子批量) |
|
||||
> | `rdb.TxPipeline()` | MULTI/EXEC 事务包装(原子批量) |
|
||||
> | `rdb.Pipeline()` | 普通 Pipeline(非原子批量,仅合并发送) |
|
||||
> | `rdb.TxPipeline()` | MULTI/EXEC 事务包装(命令串行执行) |
|
||||
>
|
||||
> **生产环境推荐用 `TxPipeline()`**,因为它保证了操作的原子性。
|
||||
> | 场景 | 选择 |
|
||||
> |------|------|
|
||||
> | 只需要减少 RTT,每条命令独立执行 | `rdb.Pipeline()` |
|
||||
> | 需要多条命令作为一个整体串行执行 | `rdb.TxPipeline()` |
|
||||
> | 需要乐观锁(WATCH + CHECK) | `rdb.TxPipeline()` |
|
||||
>
|
||||
> **限流场景**中,每个命令通常操作不同的 Key、各自独立判断——用 `Pipeline()` 即可。
|
||||
> **数据一致性要求高**(如统计计数 + 过期设置绑定)——用 `TxPipeline()`。
|
||||
|
||||
## 三、在限流中的应用
|
||||
|
||||
### 3.1 多级限流的 Pipeline 优化
|
||||
|
||||
原始方式三次独立 Redis 调用,每次都要经历完整的网络往返。Pipeline 可以将其压缩为两次往返(L1 单独 + L2/L3 合并)。
|
||||
|
||||
```go
|
||||
// 原始方式:三次独立的 Redis 调用
|
||||
// 原始方式:三次独立的 Redis 调用 → 3× RTT
|
||||
func multiLevelLimitSlow(ctx context.Context, rdb *redis.Client, ip, userID, endpoint string) error {
|
||||
if err := checkIPLimit(ctx, rdb, ip); err != nil {
|
||||
return err // L1 拦截
|
||||
return err // L1 IP 拦截
|
||||
}
|
||||
if err := checkUserLimit(ctx, rdb, userID); err != nil {
|
||||
return err // L2 拦截
|
||||
return err // L2 用户级拦截
|
||||
}
|
||||
if err := checkTokenBucket(ctx, rdb, endpoint); err != nil {
|
||||
return err // L3 拦截
|
||||
return err // L3 TokenBucket 拦截
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Pipeline 优化:L2 + L3 合并为一个批次
|
||||
// Pipeline 优化:L2 + L3 合并为一个批次 → 2× RTT
|
||||
func multiLevelLimitFast(ctx context.Context, rdb *redis.Client, ip, userID, endpoint string) error {
|
||||
// L1 单独调用(不可合并到同一个 Key space)
|
||||
// L1 单独调用(不可合并到同一个 Key space,需要串行拦截优先级)
|
||||
if err := checkIPLimit(ctx, rdb, ip); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
|
||||
// L2 + L3 用 Pipeline 打包
|
||||
pipe := rdb.TxPipeline()
|
||||
|
||||
pipe := rdb.Pipeline() // 独立判断,不需要事务语义
|
||||
|
||||
// slidingWindowLua 参数说明:KEYS[1]=rate:user:{userID}, ARGV[1]=maxCount(86400), ARGV[2]=limit(1000)
|
||||
userResult := pipe.Eval(ctx, slidingWindowLua, []string{fmt.Sprintf("rate:user:{%s}", userID)}, 86400, 1000)
|
||||
|
||||
// tokenBucketLua 参数说明:KEYS[1]=rate:tokenbucket:{endpoint}, ARGV[1]=capacity(100), ARGV[2]=refillRate(10), ARGV[3]=tokens(1), ARGV[4]=timestamp(ms)
|
||||
tokenResult := pipe.Eval(ctx, tokenBucketLua, []string{fmt.Sprintf("rate:tokenbucket:{%s}", endpoint)}, 100, 10, 1, time.Now().UnixMilli())
|
||||
|
||||
|
||||
_, err := pipe.Exec(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// 检查结果
|
||||
|
||||
// 依次检查 L2 / L3 结果
|
||||
userCount, _ := userResult.Int()
|
||||
if userCount >= 1000 {
|
||||
return ErrRateLimited
|
||||
}
|
||||
|
||||
tokenResult, _ := tokenResult.IntSlice()
|
||||
if tokenResult[0] == 0 {
|
||||
|
||||
tkArr, _ := tokenResult.IntSlice()
|
||||
if tkArr[0] == 0 {
|
||||
return ErrRateLimited
|
||||
}
|
||||
|
||||
|
||||
return nil
|
||||
}
|
||||
```
|
||||
|
||||
### 3.2 ZAdd + EXPIRE 的 Pipeline 优化
|
||||
|
||||
当滑动窗口算法不使用 Lua 脚本时(牺牲部分原子性换取代码简单),可以用 Pipeline 把"清理旧数据 → 判断数量 → 写入新数据 → 设置过期"分批执行。
|
||||
|
||||
```go
|
||||
// 不用 Lua 时的简化写法(牺牲部分原子性换取简单)
|
||||
func SlidingWindowPipeline(ctx context.Context, rdb *redis.Client, id string, windowSec int, maxLen int64) error {
|
||||
@@ -146,20 +176,22 @@ func SlidingWindowPipeline(ctx context.Context, rdb *redis.Client, id string, wi
|
||||
cutoff := now - int64(windowSec)*1000
|
||||
member := ulid.Now().String()
|
||||
|
||||
pipe := rdb.TxPipeline()
|
||||
// 第一批次:清理过期记录 + 读取当前数量
|
||||
pipe := rdb.Pipeline()
|
||||
pipe.ZRemRangeByScore(ctx, key, "-inf", strconv.FormatInt(cutoff, 10))
|
||||
pipe.ZCard(ctx, key)
|
||||
results, err := pipe.Exec(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
|
||||
currentCount := int(results[1].(*redis.IntCmd).Val())
|
||||
if currentCount >= int(maxLen) {
|
||||
return ErrRateLimited
|
||||
}
|
||||
|
||||
pipe2 := rdb.TxPipeline()
|
||||
|
||||
// 第二批次:写入新记录 + 设置 Key 过期时间
|
||||
pipe2 := rdb.Pipeline()
|
||||
pipe2.ZAdd(ctx, key, &redis.Z{Score: float64(now), Member: member})
|
||||
pipe2.Expire(ctx, key, time.Duration(windowSec)*time.Second)
|
||||
_, err = pipe2.Exec(ctx)
|
||||
@@ -167,35 +199,35 @@ func SlidingWindowPipeline(ctx context.Context, rdb *redis.Client, id string, wi
|
||||
}
|
||||
```
|
||||
|
||||
> [!question]- 这里有两批 Pipeline,能合并成一批吗?
|
||||
>
|
||||
> 想一想:`ZCard` 返回的数量决定了是否允许 `ZAdd`——这是一个**读 → 判断 → 写**的模式。Pipeline 无法根据第一条命令的结果决定是否执行第二条,所以必须分成两批。这种情况应该考虑什么方案?
|
||||
|
||||
## 四、Pipeline 的注意事项
|
||||
|
||||
### 4.1 注意事项汇总
|
||||
### 4.1 常见陷阱
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
subgraph "⚠️ Pipeline 陷阱"
|
||||
T1["不能用 Pipeline 做条件逻辑<br/>无法根据 A 的结果决定 B"] --> T2["需要判断分支时用 Lua 脚本"]
|
||||
|
||||
T3["大批量发送可能超过 Redis<br/>maxpacket 限制"] --> T4["控制单次 Pipeline 的命令数<br/>建议 50~200 条"]
|
||||
|
||||
T5["流水线中的错误不会中断后续命令"] --> T6["需要在客户端检查每个命令的返回"]
|
||||
|
||||
T7["Cluster 模式下多 Key<br/>必须在同一 slot"] --> T8["用 Hash Tag {} 保证同 slot"]
|
||||
end
|
||||
|
||||
style T2 fill:#e3f2fd
|
||||
style T4 fill:#fff3e0
|
||||
style T6 fill:#fff3e0
|
||||
style T8 fill:#e3f2fd
|
||||
flowchart LR
|
||||
A["不能用 Pipeline<br/>做条件逻辑"] --> B["读 → 判断 → 写<br/>需要 Lua 脚本"]
|
||||
C["大批量可能超过<br/>Redis 内部限制"] --> D["控制单次 50~200 条"]
|
||||
E["错误不会中断<br/>后续命令"] --> F["客户端逐个检查结果"]
|
||||
G["Cluster 多 Key<br/>必须在同一 slot"] --> H["用 Hash Tag {} 同槽"]
|
||||
|
||||
classDef warn fill:#fff3e0
|
||||
classDef info fill:#e3f2fd
|
||||
class B,H info
|
||||
class D,F warn
|
||||
```
|
||||
|
||||
### 4.2 最佳实践
|
||||
|
||||
| 场景 | 推荐方式 | 原因 |
|
||||
|------|---------|------|
|
||||
| 多条独立写入 | Pipeline (`TxPipeline`) | 减少 RTT,保持简单 |
|
||||
| 读 → 判断 → 写 | Lua 脚本 | 需要原子性和业务逻辑 |
|
||||
| 批量删除大 Key | UNLINK (逐个) | DEL 在大 Key 时会阻塞 |
|
||||
| 多条独立写入(如批量 SET) | `Pipeline()` | 只需减少 RTT,无需事务 |
|
||||
| 需要原子顺序执行(如 INCR + EXPIRE 绑定) | `TxPipeline()` | 避免被其他客户端插入 |
|
||||
| 读 → 判断 → 写 | Lua 脚本 | Pipeline 无法做条件分支 |
|
||||
| 批量删除大 Key | UNLINK(逐个) | DEL 在大 Key 时会阻塞主线程 |
|
||||
| 统计类聚合操作 | Lua 中遍历 | 避免多次往返的数据不一致 |
|
||||
|
||||
## 五、性能基准参考
|
||||
@@ -212,8 +244,13 @@ flowchart TD
|
||||
>
|
||||
> 对于限流这种高频场景(每秒数千到数万请求),即使是 1ms 的额外 RTT 也可能成为瓶颈。Pipeline 的价值在于将 N 次 RTT 压缩为 1 次。
|
||||
|
||||
> [!question]- Pipeline 是银弹吗?
|
||||
>
|
||||
> Pipeline 虽然减少了网络延迟,但它不能解决所有问题——无法在流水线中做条件分支、大批量可能触达 Redis 内部限制。什么时候该用 Pipeline、什么时候该换 Lua 脚本?答案就在上一节和下一节的对比中。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[分布式限流]] — 性能优化清单中的 Pipeline 章节
|
||||
- [[Go-Redis Lua 调用指南]] — 与 EvalSHA 配合使用的组合策略
|
||||
- [[SCAN命令]] — SCAN 遍历时可以结合 Pipeline 提高批量处理效率
|
||||
- [[MULTI/EXEC]] — Redis 事务详解
|
||||
|
||||
Reference in New Issue
Block a user