Files
cs-note/hzh/REDIS/Pipeline批量操作.md

9.6 KiB
Raw Permalink Blame History

tags, create time
tags create time
redis
pipeline
batch
performance
multi-exec
2026-06-01 10:30

Pipeline 批量操作

概述

Redis 的每次命令执行都有网络往返(RTT)开销。对于需要执行多条命令的场景,Pipeline 可以把多个命令打包成一个批次发送,大幅降低 RTT 次数。这是限流性能优化的核心手段之一(F14 考点)。

[!question]- 思考一下

假设你的应用和 Redis 在同一机房(RTT ≈ 0.5ms),每条命令平均处理时间 0.1ms。如果一次请求需要 10 条 Redis 命令,总延迟是多少?用 Pipeline 打包后呢?

一、为什么需要 Pipeline?

无 Pipeline:逐条发送

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:批量发送

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% 的网络延迟。

二、Pipeline vs MULTI/EXEC

这是最容易混淆的概念。填空题 F14 的答案是 MULTI/EXEC,但实际使用中有重要的区别。

对比表

维度 Pipeline MULTI/EXEC
原子性 ❌ 每条命令独立执行 ✅ EXEC 时整体执行
取消支持 ❌ 不能中途取消 ✅ DISCARD 取消
事务回滚 ❌ 单条失败不影响其他 ⚠️ 编译错误时不执行,运行时错误照常返回
嵌套管道 ❌ 不能在事务内用 Pipeline ❌ 不支持嵌套
Watch 支持 ❌ 无 ✅ 配合 WATCH 实现乐观锁

关键区别:原子性

[!warning]- Redis 事务的"坑"

Redis 的 MULTI/EXEC 不提供回滚语义!如果某条命令在 EXEC 时因为类型错误等运行时异常失败,其他命令仍会执行。它只保证 EXEC 后的命令不被其他客户端打断——也就是串行化,而非传统数据库的 ACID 事务。

// Pipeline:每条命令独立执行,前一条成功后面失败也照常返回
pipe := rdb.Pipeline()
pipe.Incr(ctx, "key1")          // 成功写入
pipe.Expire(ctx, "key1", 60)    // 即使这步出错,Incr 结果仍会返回
results, _ := pipe.Exec(ctx)

// TxPipeline:MULTI/EXEC 包装,命令串行执行(不被其他客户端插入)
// 注意:不是真正的原子回滚!
pipe2 := rdb.TxPipeline()       // = MULTI ... EXEC 包装
pipe2.Incr(ctx, "key1")
pipe2.Expire(ctx, "key1", 60)
results2, _ := pipe2.Exec(ctx)

如何选择:Pipeline 还是 TxPipeline?

[!tip]- Go-Redis 的 API 选择指南

API 对应行为
rdb.Pipeline() 普通 Pipeline(非原子批量,仅合并发送)
rdb.TxPipeline() MULTI/EXEC 事务包装(命令串行执行)
场景 选择
只需要减少 RTT,每条命令独立执行 rdb.Pipeline()
需要多条命令作为一个整体串行执行 rdb.TxPipeline()
需要乐观锁(WATCH + CHECK) rdb.TxPipeline()

限流场景中,每个命令通常操作不同的 Key、各自独立判断——用 Pipeline() 即可。 数据一致性要求高(如统计计数 + 过期设置绑定)——用 TxPipeline()。

三、在限流中的应用

3.1 多级限流的 Pipeline 优化

原始方式三次独立 Redis 调用,每次都要经历完整的网络往返。Pipeline 可以将其压缩为两次往返(L1 单独 + L2/L3 合并)。

// 原始方式:三次独立的 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 IP 拦截
    }
    if err := checkUserLimit(ctx, rdb, userID); err != nil {
        return err // L2 用户级拦截
    }
    if err := checkTokenBucket(ctx, rdb, endpoint); err != nil {
        return err // L3 TokenBucket 拦截
    }
    return nil
}

// Pipeline 优化:L2 + L3 合并为一个批次 → 2× RTT
func multiLevelLimitFast(ctx context.Context, rdb *redis.Client, ip, userID, endpoint string) error {
    // L1 单独调用(不可合并到同一个 Key space,需要串行拦截优先级)
    if err := checkIPLimit(ctx, rdb, ip); err != nil {
        return err
    }

    // L2 + L3 用 Pipeline 打包
    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
    }

    tkArr, _ := tokenResult.IntSlice()
    if tkArr[0] == 0 {
        return ErrRateLimited
    }

    return nil
}

3.2 ZAdd + EXPIRE 的 Pipeline 优化

当滑动窗口算法不使用 Lua 脚本时(牺牲部分原子性换取代码简单),可以用 Pipeline 把"清理旧数据 → 判断数量 → 写入新数据 → 设置过期"分批执行。

// 不用 Lua 时的简化写法(牺牲部分原子性换取简单)
func SlidingWindowPipeline(ctx context.Context, rdb *redis.Client, id string, windowSec int, maxLen int64) error {
    key := fmt.Sprintf("rate:sliding:%s", id)
    now := time.Now().UnixMilli()
    cutoff := now - int64(windowSec)*1000
    member := ulid.Now().String()

    // 第一批次:清理过期记录 + 读取当前数量
    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
    }

    // 第二批次:写入新记录 + 设置 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)
    return err
}

[!question]- 这里有两批 Pipeline,能合并成一批吗?

想一想:ZCard 返回的数量决定了是否允许 ZAdd——这是一个读 → 判断 → 写的模式。Pipeline 无法根据第一条命令的结果决定是否执行第二条,所以必须分成两批。这种情况应该考虑什么方案?

四、Pipeline 的注意事项

4.1 常见陷阱

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 最佳实践

场景 推荐方式 原因
多条独立写入(如批量 SET) Pipeline() 只需减少 RTT,无需事务
需要原子顺序执行(如 INCR + EXPIRE 绑定) TxPipeline() 避免被其他客户端插入
读 → 判断 → 写 Lua 脚本 Pipeline 无法做条件分支
批量删除大 Key UNLINK(逐个) DEL 在大 Key 时会阻塞主线程
统计类聚合操作 Lua 中遍历 避免多次往返的数据不一致

五、性能基准参考

假设网络 RTT = 1ms,本地部署(RTT ≈ 0.1ms):

命令数 无 Pipeline (RTT×N) Pipeline (1 RTT) 加速比
10 条 10ms / 1ms 1ms / 0.1ms 10x
100 条 100ms / 10ms 1ms / 0.1ms 100x
1000 条 1000ms / 100ms 1ms / 0.1ms 1000x

[!note]- 实际应用

对于限流这种高频场景(每秒数千到数万请求),即使是 1ms 的额外 RTT 也可能成为瓶颈。Pipeline 的价值在于将 N 次 RTT 压缩为 1 次。

[!question]- Pipeline 是银弹吗?

Pipeline 虽然减少了网络延迟,但它不能解决所有问题——无法在流水线中做条件分支、大批量可能触达 Redis 内部限制。什么时候该用 Pipeline、什么时候该换 Lua 脚本?答案就在上一节和下一节的对比中。

关联笔记