Skip to content

第 7 章 事务、Pipeline、Lua、Pub/Sub、Stream

学习目标:把 Redis 提供的「多命令组合武器库」一次性盘清楚——什么时候用事务、什么时候用 Pipeline、什么时候用 Lua、什么时候用消息队列。看完之后能干净利落地回答:「为什么 Redis 事务没有回滚」「Pipeline 跟 MULTI 到底差在哪」「Lua 脚本为什么是原子的」「为什么生产环境推荐 Stream 而不是 Pub/Sub」「Stream 的 Consumer Group 跟 Kafka 像不像」。


7.1 概览:本章 5 个相关但独立的特性

打开本章之前先树立一个观念——这 5 个特性虽然常被放在一起讲,但它们解决的根本问题完全不同

┌─────────────┬──────────────────────────────┬───────────────────┐
│  特性        │  解决什么问题                  │  关键词            │
├─────────────┼──────────────────────────────┼───────────────────┤
│ 事务         │  多条命令打包「一次性按顺序执行」│ MULTI/EXEC/WATCH  │
│ Pipeline    │  减少 N 次网络往返(RTT)       │ 批量发/批量收      │
│ Lua 脚本     │  把「读改写」逻辑下沉到服务端    │ EVAL/原子/共享解释器│
│ Pub/Sub     │  广播通知,订阅者实时收到         │ PUBLISH/SUBSCRIBE │
│ Stream      │  真正的消息队列:持久化 + 消费组   │ XADD/XREADGROUP   │
└─────────────┴──────────────────────────────┴───────────────────┘

🍱 生活类比

  • 事务 = 你去银行办业务:取号 → 排队 → 在窗口一次办完所有事情,中途不会被插队。
  • Pipeline = 你去快递柜取 5 个包裹:把 5 个取件码一次性输入,而不是排队 5 次。
  • Lua 脚本 = 你给银行写一份办事流程纸条,让柜员按流程走完,期间你在外面等。
  • Pub/Sub = 公园广播:所有正在听的人都听见了,没在听的人错过就错过。
  • Stream = 取号机+叫号屏:每个号码都被记下来,叫到了你也能后悔重新查询。

下面挨个拆解。


7.2 Redis 事务(MULTI / EXEC / DISCARD / WATCH)

7.2.1 是什么 + 生活类比

Redis 的事务就是把多条命令打包,一次性按顺序发给服务端执行,中途不会被其他客户端的命令打断

🏦 生活类比:你到银行柜台办「转账 + 改预留手机号 + 打印流水」3 件事。柜员先把 3 张单子收齐(MULTI 阶段:命令排队),然后挂出「业务处理中」的牌子(EXEC 阶段:单线程一口气执行),办完才放下一个客户进来。

7.2.2 4 个核心命令

bash
MULTI                  # 开启事务(开始排队)
SET balance 100        # → 返回 QUEUED,命令进入队列
INCRBY balance 50      # → QUEUED
GET balance            # → QUEUED
EXEC                   # 一次性执行队列里所有命令,返回结果数组

# 或者反悔:
MULTI
SET k1 v1
DISCARD                # 放弃整个事务,清空队列

实测一下:

bash
127.0.0.1:6379> MULTI
OK
127.0.0.1:6379(TX)> SET k1 hello
QUEUED
127.0.0.1:6379(TX)> SET k2 world
QUEUED
127.0.0.1:6379(TX)> GET k1
QUEUED
127.0.0.1:6379(TX)> EXEC
1) OK
2) OK
3) "hello"

7.2.3 排队过程到底发生了什么

客户端 A                        服务端
   │   MULTI                       │
   ├──────────────────────────────►│  → flag CLIENT_MULTI = 1
   │   OK                           │     新建 multiState 队列
   │◄──────────────────────────────┤
   │                                │
   │   SET k1 v1                    │
   ├──────────────────────────────►│  → 不执行,仅做语法检查
   │   QUEUED                       │     检查通过则塞进队列
   │◄──────────────────────────────┤
   │                                │
   │   SET k2 v2                    │
   ├──────────────────────────────►│
   │   QUEUED                       │
   │◄──────────────────────────────┤
   │                                │
   │   EXEC                         │
   ├──────────────────────────────►│  → 串行 pop 队列里每条命令
   │                                │     一条条调用执行函数
   │                                │     ⚠️ 期间不会被其他客户端打断
   │                                │       (Redis 主线程串行)
   │   1) OK / 2) OK                │
   │◄──────────────────────────────┤

7.2.4 WATCH:实现乐观锁

光有 MULTI/EXEC 还不够——如果 EXEC 之前别的客户端把你要操作的 Key 改了怎么办?答案是 WATCH

🪜 生活类比:你想从冰箱拿最后一瓶可乐,先在瓶子上贴张便利贴写「我在看」(WATCH),然后去拿杯子(MULTI 排队),回来之前如果别人撕掉便利贴动过瓶子,你就放弃喝(EXEC 返回 nil)。

bash
WATCH balance                # 监视 balance
val = GET balance            # 读到 100
MULTI
SET balance (val - 30)       # 想扣 30 块
EXEC                         # 如果 balance 在 WATCH 之后被改过,整个事务作废,返回 (nil)

WATCH 的实现原理

db.watched_keys: { key → list[clients] }   ← 每个 Key 维护「谁在看」

任何写命令在执行时(如 SET k1 ...):
  1. 查 watched_keys 看 k1 有没有客户端在 watch
  2. 如果有,给那些客户端打上 CLIENT_DIRTY_CAS 标记

EXEC 执行前:
  1. 检查自己有没有 DIRTY_CAS 标记
  2. 有 → 直接放弃事务,返回 nil
  3. 无 → 正常执行队列

这就是乐观锁:不主动加锁阻塞别人,而是在提交时检查「我看的东西有没有被动过」,被动了就回滚重试。

7.2.5 完整的乐观锁伪代码

python
while True:
    pipe.watch("balance")
    cur = int(pipe.get("balance"))
    if cur < 30:
        pipe.unwatch()
        raise InsufficientFund()
    pipe.multi()
    pipe.set("balance", cur - 30)
    try:
        pipe.execute()        # 失败抛 WatchError
        break
    except WatchError:
        continue              # 重试

7.2.6 ⚠️ 为什么 Redis 事务不是真事务

经典面试题。Redis 事务跟数据库事务的对比:

┌─────────────┬──────────────────┬────────────────────┐
│  ACID 属性   │  数据库事务         │  Redis 事务         │
├─────────────┼──────────────────┼────────────────────┤
│ A 原子性     │ ✅ 失败全部回滚      │ ⚠️ 入队错误整体不执行│
│             │                   │    执行错误其他命令   │
│             │                   │    照常执行(无回滚) │
│ C 一致性     │ ✅ 约束/外键保证     │ ✅ 单线程执行        │
│ I 隔离性     │ ✅ 锁/MVCC          │ ✅ 单线程天然隔离    │
│ D 持久性     │ ✅ WAL              │ ⚠️ 取决于 AOF 配置  │
└─────────────┴──────────────────┴────────────────────┘

最关键的反差:Redis 事务没有回滚

7.2.7 入队错误 vs 执行时错误

这是两类完全不同的错误,处理方式也不一样:

─── 入队错误(语法错) ──────────────────────────
MULTI
SETT k1 v1            ← 命令名拼错
EXEC

→ 服务端在 EXEC 时直接拒绝整个事务
  (error) EXECABORT Transaction discarded because of previous errors.
  ✅ 这种情况是「全有或全无」,相当于回滚

─── 执行时错误(类型错) ────────────────────────
MULTI
SET counter "abc"
INCR counter          ← 字符串 "abc" 不能 INCR
SET k1 v1
EXEC

→ 1) OK
   2) (error) value is not an integer
   3) OK
   ⚠️ 其他命令照常执行!中间的错误命令不会让前后回滚

为什么不做回滚? 作者 antirez 给出的理由:

  1. Redis 命令大部分错误是编程错误(类型用错、参数用错),生产环境不该出现,不该为了少数 bug 让所有事务付出回滚的复杂度代价。
  2. 没有回滚,Redis 事务的实现可以保持简单且更快
  3. AOF/RDB 持久化追加单条命令即可,无需 undo log。

7.2.8 事务的局限

局限说明
没有真正回滚中间命令报错,前后命令照常执行
不能根据上一条结果改下一条队列里命令的参数必须是确定值,因为还没执行
WATCH 在 Cluster 下要求所有 Key 同槽跨槽事务不支持,需要 hashtag
大量小事务不如 Lua 脚本 + Pipeline 高效

7.3 Pipeline 批处理

7.3.1 是什么 + 生活类比

Pipeline 是纯客户端侧的优化——把多条命令攒一批一次性发给服务端,再一次性收回所有响应,从而消除大量网络往返(RTT)开销

📦 生活类比

  • 不用 Pipeline = 你在快递柜取 5 个包裹:输入第 1 个码 → 等开柜门 → 拿出 → 关门 → 输入第 2 个码 → ……(来回 5 次)
  • 用 Pipeline = 把 5 个取件码一次性输入,柜门同时打开,5 个包裹一次取完。
┌──────────────────────────────────────────────────────────┐
│ 不用 Pipeline(每条命令一次 RTT)                          │
├──────────────────────────────────────────────────────────┤
│ Client → SET k1 v1 → Server   ⏳ RTT                      │
│ Client ← OK ←──────── Server                              │
│ Client → SET k2 v2 → Server   ⏳ RTT                      │
│ Client ← OK ←──────── Server                              │
│ Client → SET k3 v3 → Server   ⏳ RTT                      │
│ Client ← OK ←──────── Server                              │
│   总耗时 ≈ N × RTT + N × 服务端处理                        │
└──────────────────────────────────────────────────────────┘

┌──────────────────────────────────────────────────────────┐
│ 用 Pipeline                                              │
├──────────────────────────────────────────────────────────┤
│ Client → [SET k1 v1] [SET k2 v2] [SET k3 v3] → Server    │
│                                       ⏳ 1 个 RTT          │
│ Client ← [OK] [OK] [OK] ←──────────── Server              │
│   总耗时 ≈ 1 × RTT + N × 服务端处理                        │
└──────────────────────────────────────────────────────────┘

7.3.2 Pipeline ≠ 事务!

这是面试官最爱挖的坑。两者都「一批发出去」,但语义完全不同:

┌──────────────────┬────────────────┬────────────────────┐
│  对比项           │  Pipeline       │  MULTI/EXEC 事务    │
├──────────────────┼────────────────┼────────────────────┤
│ 客户端攒命令       │ ✅              │ ✅                 │
│ 一次发给服务端     │ ✅              │ ✅                 │
│ 期间能被别人打断   │ ⚠️ 能           │ ❌ 不能(原子)     │
│ 失败影响后续       │ ⚠️ 不影响        │ ⚠️ 不影响(也无回滚)│
│ 主要目的           │ 省 RTT          │ 原子 + 顺序        │
│ 服务端额外开销     │ 0               │ 维护 multi 队列    │
└──────────────────┴────────────────┴────────────────────┘

要点:

  • Pipeline 不保证原子性:服务端在处理你的批命令时,如果中间收到其他客户端的命令也会去处理(除非你在 Pipeline 里嵌入 MULTI/EXEC)。
  • Pipeline 只是网络层优化:你的客户端不等响应就连续发,服务端依次处理,响应攒着一起返回(其实是流式返回的,TCP 缓冲区凑批)。

💡 实际工程中常见用法是 Pipeline + MULTI/EXEC 套娃:先用 Pipeline 减少 RTT,再用 MULTI/EXEC 保证原子性。redis-pypipeline(transaction=True)(默认值)就是这个意思。

7.3.3 Pipeline 实测:1 万次 SET 提速 ~50x

跑一下 code/02_pipeline_benchmark.py 你会看到类似输出:

方式               耗时         qps
──────────────────────────────────────
单条 SET           5.20 s       1923
Pipeline (100/批)  0.10 s     100000
MULTI/EXEC         0.18 s      55555

为什么差距这么大?因为本机 RTT 即便只有 0.5ms,1 万次也是 5 秒;Pipeline 把 RTT 摊到几十次。网络越远(跨机房 / 上云),Pipeline 提速越夸张

7.3.4 Pipeline 注意事项

  1. 批大小别太大:单批超过几 MB 会撑爆服务端 query buffer,建议 100 ~ 1000 一批。
  2. 响应必须按顺序读完:客户端拿到响应数组的顺序就是发出命令的顺序。
  3. Cluster 下要 KEY 同槽:跨槽 Key 无法在同一个 Pipeline 里发(要拆分到不同节点)。
  4. 不影响其他客户端:Pipeline 本身不阻塞,只是连续发命令。

7.4 Lua 脚本

7.4.1 是什么 + 生活类比

Lua 脚本让你把一段读改写逻辑(甚至带 if/else / 循环)一次性发到服务端执行,期间不会被其他命令打断——本质是在服务端攒了个原子化的小函数

🍔 生活类比:你去快餐店点餐,每加一个配料就跑去问老板「加这个能不能」「加那个能不能」实在太麻烦。你不如写好一张「点餐流程纸条」(Lua 脚本):先看库存,有就扣 1 份;没有就报错。把纸条递给收银员,他一次办完。

7.4.2 EVAL / EVALSHA 用法

bash
# 直接执行脚本
EVAL "return redis.call('GET', KEYS[1])" 1 mykey
                    └────脚本────┘   └─key─┘
                                  KEYS 数量

# KEYS 与 ARGV 的区别:
EVAL "return {KEYS[1], ARGV[1]}" 1 myKey myArg
# → 1) "myKey"
#   2) "myArg"

# 缓存脚本,下次用 SHA1 哈希调用,省传输
SCRIPT LOAD "return 'hi'"
# → "55482b1e3d2bb5f4e..."(脚本的 SHA1)
EVALSHA 55482b1e3d2bb5f4e... 0

关键参数解释

  • EVAL <script> <numkeys> [key ...] [arg ...]
  • numkeys:告诉 Redis 后面紧跟的 N 个参数是 Key
  • KEYS[i] / ARGV[i]:在 Lua 脚本里访问,Lua 数组下标从 1 开始

7.4.3 为什么必须显式声明 KEYS

不是 Redis 强制的语法要求,而是集群路由的需要

单机版:传不传 KEYS 都能跑(KEYS[1] 也可以塞 ARGV)
Cluster: 必须显式 KEYS!
        Proxy/客户端要根据 KEYS 计算 slot,
        所有 KEYS 必须落在同一个 slot
        否则报 CROSSSLOT 错误

7.4.4 Lua 脚本的原子性

核心要点:脚本执行期间,主线程不会处理其他客户端的命令。 所以你的「读 → 判断 → 写」无论多复杂,都是原子的。

T0: Client A:EVAL <script> ...
T0+ε: 服务端开始执行脚本第一行
       ⚡ 此时 Client B 发来命令 → 排队等待
T1: 脚本最后一行执行完
T1+ε: 服务端开始处理 Client B 的命令

这就是为什么「库存扣减」这种竞争场景特别适合 Lua:

lua
-- KEYS[1] 是库存 key,ARGV[1] 是要扣的数量
local stock = tonumber(redis.call('GET', KEYS[1]) or "0")
local need  = tonumber(ARGV[1])
if stock < need then
  return -1                        -- 库存不足
end
redis.call('DECRBY', KEYS[1], need)
return stock - need                -- 返回剩余

整个过程没有 GET → 业务判断 → DECR 之间的窗口,杜绝超卖。

7.4.5 经典示例:库存扣减

完整代码见 07_transaction_pipeline/code/03_lua_stock.py

python
LUA_DECR_STOCK = """
local stock = tonumber(redis.call('GET', KEYS[1]) or '0')
local need  = tonumber(ARGV[1])
if stock < need then
    return -1
end
redis.call('DECRBY', KEYS[1], need)
return stock - need
"""

r.set("stock:item:1001", 100)
remain = r.eval(LUA_DECR_STOCK, 1, "stock:item:1001", 1)
print(remain)  # → 99(剩余库存)

实战收益:1000 个并发线程同时扣,最终剩余永远精确等于 100 - 1000(如果初始足);纯 Python 客户端做 get + 判断 + decr 一定会超卖。

7.4.6 Lua 脚本的限制与坑

限制解释应对
执行期间阻塞主线程脚本太长,所有客户端都得等控制脚本执行 < 几毫秒
超过 lua-time-limit 默认 5s慢日志告警,但不会自动停SCRIPT KILL 强杀(写过的不杀)
不能用全局变量Redis 禁止 x = 1,要 local x = 1必须 local
随机/时间不确定math.random / 系统时间不能用,影响主从复制一致性用 ARGV 传入
Cluster 必须同槽否则路由失败用 hashtag {}
脚本缓存清空风险SCRIPT FLUSH 后 EVALSHA 找不到客户端要做 NOSCRIPT 重试 → EVAL 重传

🆕 Redis 7 引入 Functions:长期持久化的命名脚本(FUNCTION LOAD),可以替代「客户端缓存 + EVALSHA」的复杂模式。


7.5 发布订阅 Pub/Sub

7.5.1 是什么 + 生活类比

Pub/Sub 是最古老的消息通道:发布者把消息扔到「频道」上,所有正在订阅这个频道的客户端当场实时收到。

📻 生活类比:公园里的广播喇叭。广播员(PUBLISH)一喊「XX 找妈妈」,正在公园的人(SUBSCRIBE)都能听到;但 5 分钟之前进来又走了的人,错过就错过了——广播不会被录下来

7.5.2 命令

bash
# 终端 1:订阅
SUBSCRIBE news weather              # 订阅 news 和 weather 频道
# → Reading messages... (press Ctrl-C to quit)

# 终端 2:发布
PUBLISH news "今天 Redis 7.4 发布"   # → (integer) 1   返回有 1 个订阅者收到

# 终端 1 立即收到:
# 1) "message"
# 2) "news"
# 3) "今天 Redis 7.4 发布"

# 模式订阅(通配符)
PSUBSCRIBE news.*                    # 订阅所有 news.开头的频道
PUBLISH news.tech "..."              # 也能收到

# 退订
UNSUBSCRIBE news
PUNSUBSCRIBE news.*

7.5.3 ⚠️ Pub/Sub 的致命缺陷

┌─────────────────────────────────────────────────────────┐
│ Pub/Sub 不能用作「严肃」消息队列!                          │
├─────────────────────────────────────────────────────────┤
│ ❌ 不持久化:消息发完即弃                                  │
│ ❌ 订阅者掉线 → 那段时间的消息全丢                          │
│ ❌ 服务端不留 backlog,消费者慢点消费会被踢                 │
│ ❌ 没有 ACK,没有重试                                     │
│ ❌ 单条消息所有订阅者都收到(不是「一条只给一个 worker 处理」)│
│ ❌ Cluster 下 PUBLISH 默认全节点广播(流量爆炸)            │
└─────────────────────────────────────────────────────────┘

那它适合什么?广播型实时通知

  • 配置变更通知(所有应用实例刷新本地缓存)
  • 多服务器之间的「失效广播」
  • 在线聊天室的房间广播
  • Redis 自身的 __keyspace@__ keyspace notifications

7.5.4 内部数据结构

server.pubsub_channels:
   "news"     → list[Client A, Client C]
   "weather"  → list[Client B, Client C]

server.pubsub_patterns:
   "news.*"   → list[Client D]

PUBLISH 时:
   1. 查 pubsub_channels["news"] → 把消息推给 A、C
   2. 遍历 pubsub_patterns,匹配的也推

完全是内存中的链表——发完就丢,没有任何持久化结构。


7.6 Stream(5.0+):真正的 Redis 消息队列

7.6.1 是什么 + 生活类比

Stream 是 Redis 5.0 引入的append-only 日志型数据结构,专门为消息队列场景设计。可以把它理解为 Redis 自带的「迷你 Kafka」。

🎟️ 生活类比:医院取号机。

  • 每个号码(消息 ID)单调递增,记录在大屏幕上不会消失(持久化)。
  • 多个窗口(Consumer Group)协作叫号,1 号窗叫完不会被 2 号窗重复叫(消费分摊)。
  • 病人没到,号牌一直留着可以重叫(PEL 待确认列表)。
  • 病人看完拿了发票回来销号(XACK)。

7.6.2 核心命令速览

bash
# 生产消息(XADD:自动生成 ID = 时间戳-序号)
XADD orders * order_id 1001 user_id 88 amount 199
# → "1713345678901-0"

XLEN orders                                   # 长度
XRANGE orders - +                             # 全部消息
XRANGE orders 1713345678901-0 +               # 从某 ID 之后

# ── 简单消费(XREAD,类似 List 的 BLPOP) ──
XREAD COUNT 2 BLOCK 5000 STREAMS orders 0     # 从 0 开始读 2 条,阻塞 5 秒

# ── 消费组消费(XREADGROUP,重头戏) ──
XGROUP CREATE orders payproc $ MKSTREAM       # 创建消费组 payproc,从最新开始
XREADGROUP GROUP payproc worker-1 COUNT 1 BLOCK 5000 STREAMS orders >
                                              # > 表示「未传递给本组的新消息」

XACK orders payproc 1713345678901-0           # 确认消息处理完成(出 PEL)

XPENDING orders payproc                        # 查看本组待确认消息

# 死信处理:长时间没 ACK 的转给别的 worker
XCLAIM orders payproc worker-2 60000 1713345678901-0

7.6.3 消息 ID 结构

1713345678901-0
└─────┬─────┘ └┬┘
   毫秒时间戳    同毫秒内的序号
   (单调递增)    (从 0 开始)

ID 永远递增:即使时间倒退(NTP 校准),Redis 也会自动让序号继续累加保单调性。

7.6.4 Consumer Group:核心数据结构

Stream 结构:
┌─────────────────────────────────────────────────────────────┐
│ orders                                                      │
│  ├── 消息 1713345678901-0  { order_id: 1001, ... }          │
│  ├── 消息 1713345678902-0  { order_id: 1002, ... }          │
│  ├── 消息 1713345678905-0  { order_id: 1003, ... }          │
│  └── 消息 1713345678910-0  { order_id: 1004, ... }          │
│                                                             │
│  Consumer Group: "payproc"                                  │
│  ├── last_delivered_id: 1713345678910-0                     │
│  │                                                          │
│  ├── Consumer "worker-1"                                    │
│  │     PEL(Pending Entries List):                         │
│  │       1713345678901-0  delivered=2,  idle=15s            │
│  │                                                          │
│  └── Consumer "worker-2"                                    │
│        PEL:                                                 │
│          1713345678905-0  delivered=1,  idle=2s             │
└─────────────────────────────────────────────────────────────┘

消费分摊语义:同一个 Group 内,每条消息只会传递给一个 Consumer——这就是 Kafka 的 Consumer Group 模型,多个 worker 协作消费。

PEL(Pending Entries List):每个 Consumer 维护自己「领走但还没 ACK」的消息列表。如果 worker 崩溃,下次它启动用 XREADGROUP ... STREAMS orders 0(注意是 0 而非 >)就能拿回之前没处理完的消息。

7.6.5 一次完整的生产 + 消费流程

┌────────────┐      XADD            ┌──────────────────────┐
│ Producer   │ ───────────────────► │ Stream: orders       │
└────────────┘                      │  msg-1, msg-2, msg-3 │
                                    └──────────┬───────────┘

              XREADGROUP ... > (新消息) ◄──────┤
              ┌────────────┐                   │
              │ worker-1   │ ← msg-1, msg-2    │  Group: payproc
              └─────┬──────┘                   │  ├─ last_delivered_id ↑
                    │ 处理完...                 │  └─ PEL[worker-1] = {msg-1, msg-2}
                    │ XACK msg-1                │
                    ▼                            │
              ┌────────────┐                   │
              │ worker-2   │ ← msg-3 (并行)    │  PEL[worker-2] = {msg-3}
              └────────────┘                   │

              超时 → XCLAIM 给 worker-3 ◄──────┘

7.6.6 Stream vs List vs Pub/Sub 全方位对比

┌───────────────────┬──────────┬───────────┬───────────┐
│  能力             │ Pub/Sub  │ List      │ Stream    │
├───────────────────┼──────────┼───────────┼───────────┤
│ 消息持久化         │ ❌       │ ✅        │ ✅        │
│ 消费者掉线后能补    │ ❌       │ ❌(弹了就没)│ ✅        │
│ 多消费者扇出广播    │ ✅(每人都收)│ ❌       │ ✅(多 Group)│
│ 多消费者协作分摊    │ ❌       │ ✅(BLPOP 抢)│ ✅(同 Group)│
│ ACK 机制          │ ❌       │ ❌        │ ✅        │
│ 历史回放           │ ❌       │ ❌        │ ✅(XRANGE)│
│ 消息 ID 递增       │ ❌       │ ❌        │ ✅        │
│ 阻塞读             │ -        │ BLPOP     │ XREAD BLOCK│
│ 内存占用           │ 极小     │ 中         │ 中(可裁剪)│
└───────────────────┴──────────┴───────────┴───────────┘

7.6.7 长度控制:MAXLEN

防止 Stream 无限增长撑爆内存:

bash
XADD orders MAXLEN 10000 * order_id 1001    # 最多保留 10000 条,老的丢弃
XADD orders MAXLEN ~ 10000 * order_id 1001  # ~ 表示近似裁剪,性能更好
XTRIM orders MAXLEN 10000
XTRIM orders MINID 1713345600000-0          # 按 ID 裁剪(早于该 ID 的全删)

7.7 实操:跑一遍配套代码

实战代码见 07_transaction_pipeline/code/

文件演示内容
01_transaction_watch.py用 WATCH 实现「乐观锁转账」,并发冲突时 EXEC 返回 nil 并重试
02_pipeline_benchmark.py1 万次 SET 三种方式对比:单条 / Pipeline / 事务
03_lua_stock.pyLua 库存扣减:1000 并发也不会超卖
04_stream_consumer_group.py完整生产者 + 2 个 worker 的 Consumer Group + ACK + XPENDING

浏览器演示见 07_transaction_pipeline/demo.html

  • ① 事务 + WATCH 演示(A、B 双客户端动画)
  • ② Pipeline vs 单条 vs 事务(三栏 RTT 对比动画)
  • ③ Lua 脚本沙盒(textarea 输 Lua + KEYS/ARGV 模拟 EVAL)
  • ④ Stream 消费组动画(Producer / 2 Consumer / PEL / ACK)

7.8 本章小结

┌────────────────────────────────────────────────────────┐
│                     本章核心要点                         │
├────────────────────────────────────────────────────────┤
│                                                        │
│  ① 事务 = 排队 + 一次执行 + 单线程隔离                   │
│     ⚠️ 不支持回滚,靠 WATCH 实现乐观锁                   │
│                                                        │
│  ② Pipeline = 客户端攒命令省 RTT                        │
│     ⚠️ 不是事务,没有原子性                              │
│     ⚠️ 与 MULTI/EXEC 可以套娃                           │
│                                                        │
│  ③ Lua 脚本 = 服务端原子化「读改写」                      │
│     ⚠️ 阻塞主线程,控制脚本执行时长                       │
│     ⚠️ Cluster 下 KEYS 必须同槽                         │
│                                                        │
│  ④ Pub/Sub = 实时广播,无持久化                          │
│     适合通知场景,不适合消息队列                          │
│                                                        │
│  ⑤ Stream = 真正的 Redis 消息队列                       │
│     XADD / XREADGROUP / XACK / PEL                    │
│     支持持久化、消费组、回放、死信                        │
│                                                        │
│  选择哲学:                                             │
│    要原子性 → Lua > 事务                                │
│    要省 RTT → Pipeline                                 │
│    要广播  → Pub/Sub                                   │
│    要队列  → Stream                                    │
│                                                        │
└────────────────────────────────────────────────────────┘

7.9 面试高频题

Q1:Redis 事务支持回滚吗?为什么?

考察点:对 Redis 事务 ACID 实现的理解。

标准答案

不支持回滚。Redis 事务的真实行为是:

  1. 入队阶段语法错 → 整个事务在 EXEC 时被拒绝(EXECABORT),算是「全有或全无」。
  2. 执行阶段类型错(如对字符串 INCR) → 该命令报错,前后命令照常执行不会回滚

为什么这样设计? Redis 作者 antirez 的理由:

  • 大部分执行错误是程序 bug(用错类型/参数),生产应该测试时就发现,而不是靠回滚兜底。
  • 不做回滚,事务实现简单、AOF 持久化无需 undo log、整体更快。
  • 「快」是 Redis 第一原则,为少数错误情况付出复杂度不划算。

加分项:提到 ACID 中 Redis 事务保证的是 C(命令一致性)+ I(单线程隔离),A 是「弱原子」,D 取决于 AOF 配置。


Q2:WATCH 是怎么实现乐观锁的?

考察点:CAS 机制理解 + 源码思想。

标准答案

WATCH 是基于 CAS(Compare-And-Swap)的乐观锁实现。流程:

  1. 客户端 WATCH key1 key2 → 服务端在 db.watched_keys 记录 {key → 客户端列表}
  2. 任何其他客户端的写命令在执行时,会查 watched_keys,把所有 watch 了该 key 的客户端打上 CLIENT_DIRTY_CAS 标记。
  3. 客户端 MULTI命令排队EXEC
  4. EXEC 执行前,检查自己有没有 DIRTY_CAS
    • 有 → 直接放弃事务,返回 nil
    • 无 → 正常执行队列。
python
while True:
    pipe.watch("balance")
    cur = int(pipe.get("balance"))
    pipe.multi()
    pipe.set("balance", cur - 30)
    try:
        pipe.execute()
        break
    except WatchError:
        continue   # 重试

加分项:提到 WATCH 的范围是「整个 Key 被改」就触发,不论改的是什么字段或值;并且 EXEC / DISCARD / 客户端断开后所有 WATCH 自动清空。


Q3:Pipeline 和事务的区别?

考察点:基础概念辨析,最爱被挖坑。

标准答案

维度Pipeline事务(MULTI/EXEC)
主要目的省网络 RTT保证原子+顺序执行
实现层纯客户端攒命令服务端有 multiState 队列
原子性❌ 没有✅ 单线程串行执行不打断
期间能否被别人插队✅ 能❌ 不能
命令出错处理各自独立入队错全失败 / 执行错继续
服务端开销0维护队列

关键澄清:两者不是二选一,可以组合使用。redis-pypipe = r.pipeline(transaction=True) 默认就是「Pipeline 攒命令 + MULTI/EXEC 包起来」——同时享受省 RTT 和原子性。

加分项:举一个实际场景说明区别——「批量给 1 万个 Key 设过期时间」用 Pipeline 就够(不要求原子);「秒杀扣库存 + 加流水」必须用事务/Lua(要求原子)。


Q4:为什么用 Lua 脚本能实现「读改写」原子性?

考察点:Lua 与 Redis 主线程的交互。

标准答案

Redis 在执行 Lua 脚本时,会把脚本当成一条复合命令对待——主线程从开始解释脚本第一行到执行完最后一行的整个过程中,不会处理任何其他客户端的命令(其他客户端的请求会在事件循环中排队等待)。

所以脚本里的「GET → 业务判断 → DECR」这一连串动作,对外部来说是瞬时完成的,不存在 GET 之后、DECR 之前的「窗口期」。

lua
local stock = tonumber(redis.call('GET', KEYS[1]))
if stock <= 0 then return -1 end
redis.call('DECR', KEYS[1])
return stock - 1

并发 1000 次调用上面的脚本,绝不会超卖

对比

python
# ❌ 客户端做这事会超卖
stock = int(r.get(key))
if stock > 0:
    r.decr(key)        # 中间隔了网络 RTT,别的客户端也在 GET → DECR

加分项

  • Lua 脚本本质是 服务端版的 MULTI/EXEC + 条件分支,但更强大(支持 if/else / 循环 / 调用任意命令)。
  • 注意脚本不能太长,否则阻塞主线程导致所有客户端等待,这就是 lua-time-limit(默认 5 秒)和 SCRIPT KILL 存在的原因。

Q5:Pub/Sub 为什么不能用作消息队列?Stream 解决了什么?

考察点:消息队列语义 + Redis 内存模型。

标准答案

Pub/Sub 的硬伤

  1. 不持久化:消息发出去就被立即推给在线订阅者,没有 backlog,发完即弃。
  2. 订阅者掉线丢消息:那段时间发的消息全错过,重连后看不到历史。
  3. 无 ACK 无重试:消费者收完没处理或处理失败,没机制重投。
  4. 慢消费者会被踢:服务端 output buffer 满 → 强制断开(client-output-buffer-limit pubsub)。
  5. 扇出语义固定:所有订阅者都收到同一条,做不了「多 worker 协作分摊」。

Stream 的应对

缺陷Stream 怎么解决
不持久化✅ 消息存在 Stream 里,AOF/RDB 一并落盘
掉线丢消息✅ 重连后用 XREADGROUP STREAMS x 0 拿回 PEL 里没 ACK 的
无 ACK✅ XACK 确认机制 + PEL 待确认列表
协作消费✅ Consumer Group:同 Group 内消息分摊,多 Group 之间扇出广播
死信处理✅ XPENDING 看积压 + XCLAIM 转给其他 worker
历史回放✅ XRANGE 任意时间段查询

加分项:提到 Redis 5.0 之前生产消息队列推荐方案是 List + BLPOP,但缺乏 ACK 和多消费者协作;Stream 是 antirez 借鉴 Kafka 设计的产物。


Q6:Stream 的 Consumer Group 与 Kafka 的 Consumer Group 有什么异同?

考察点:对消息队列模型的理解,区分中间件特性。

标准答案

相同点

  1. 消费分摊:同 Group 内每条消息只投递给一个 Consumer,多 worker 协作。
  2. 多 Group 扇出:不同 Group 之间消息隔离,每个 Group 都能完整消费。
  3. 位点(offset)持久化:消费进度记录在服务端,重启不丢。
  4. 支持回放:可以从指定位置重新读。

不同点

维度Redis StreamKafka
存储模型单机内存(可持久化)分布式磁盘日志(segment + index)
分区机制无原生分区,靠多 Stream 模拟Topic → Partition,原生分区
消费分摊单位单条消息(worker 抢)Partition(Partition 绑定 Consumer)
顺序性单 Stream 全局有序,Group 内不保证单 worker 顺序Partition 内有序
ACK 粒度逐条 XACKoffset 提交(一批一起)
消息重投XCLAIM 把 PEL 里的转给别人rebalance + 从上次 offset 重读
吞吐数十万 qps(单实例)百万 qps(集群)
持久化保证取决于 AOF 配置(fsync)ISR + acks=all
Pull vs PushPull(XREADGROUP 主动拉)Pull(poll)

核心差异

  • Redis Stream 适合中小规模、低延迟的「副业级队列」,比如订单异步通知、操作流水。
  • Kafka 适合海量数据流处理、跨数据中心、长期存储,是基础设施级的中间件。

加分项:提到 Stream 的 ACK 是逐条的(精确),但意味着高并发下 XACK 也是要走主线程的开销;Kafka 的 offset 批提交效率更高但精度低。


📌 下一章预告:第 8 章我们进入「分布式」领域——主从复制(Replication)。从全量同步、增量同步、复制积压缓冲区到 PSYNC 协议,把数据怎么从主节点流到从节点的全过程拆开看。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
"""
Ch7 配套代码 1 / 4 —— 事务 + WATCH 实现乐观锁转账

演示:
  1. 简单 MULTI/EXEC 的命令排队 + 一次执行
  2. WATCH 实现「读改写」乐观锁
  3. 故意制造并发冲突,看 EXEC 返回 nil 触发 WatchError 后自动重试
"""

import threading
import time
import redis

POOL = redis.ConnectionPool(host="127.0.0.1", port=6379, decode_responses=True)


def section(title: str) -> None:
    print("\n" + "=" * 62)
    print(title)
    print("=" * 62)


def demo_basic_transaction() -> None:
    section("Demo 1: 基础 MULTI/EXEC —— 命令排队 + 一次执行")
    r = redis.Redis(connection_pool=POOL)
    r.set("balance", 100)

    pipe = r.pipeline(transaction=True)
    pipe.incrby("balance", 50)
    pipe.decrby("balance", 30)
    pipe.get("balance")
    results = pipe.execute()

    print(f"  事务执行返回 = {results}")
    print(f"  balance 最终 = {r.get('balance')}  (期望 120)")
    r.delete("balance")


def demo_exec_runtime_error() -> None:
    section("Demo 2: 执行时错误 —— 其他命令依然执行(无回滚!)")
    r = redis.Redis(connection_pool=POOL)
    r.set("counter", "abc")
    r.delete("k1", "k2")

    pipe = r.pipeline(transaction=True)
    pipe.set("k1", "v1")
    pipe.incr("counter")
    pipe.set("k2", "v2")
    try:
        results = pipe.execute(raise_on_error=False)
    except Exception as e:
        results = [e]

    for i, res in enumerate(results, 1):
        print(f"  cmd{i} -> {res!r}")
    print(f"\n  k1 = {r.get('k1')}  (前面命令照常执行)")
    print(f"  k2 = {r.get('k2')}  (后面命令也照常执行)")
    print("  💡 Redis 事务不支持回滚,中间出错前后命令都生效")
    r.delete("counter", "k1", "k2")


def transfer_with_watch(from_acc: str, to_acc: str, amount: int, label: str) -> bool:
    """用 WATCH 实现乐观锁转账。返回 True = 成功,False = 余额不足。"""
    r = redis.Redis(connection_pool=POOL)
    retries = 0
    while True:
        try:
            with r.pipeline() as pipe:
                pipe.watch(from_acc, to_acc)

                from_bal = int(pipe.get(from_acc) or 0)
                to_bal = int(pipe.get(to_acc) or 0)

                if from_bal < amount:
                    pipe.unwatch()
                    print(f"  [{label}] ❌ 余额不足 ({from_bal} < {amount})")
                    return False

                pipe.multi()
                pipe.set(from_acc, from_bal - amount)
                pipe.set(to_acc, to_bal + amount)
                pipe.execute()
                print(f"  [{label}] ✅ 转账成功 (重试 {retries} 次)")
                return True
        except redis.WatchError:
            retries += 1
            time.sleep(0.001)
            if retries > 50:
                print(f"  [{label}] 🚫 重试过多放弃")
                return False


def demo_watch_optimistic_lock() -> None:
    section("Demo 3: WATCH 乐观锁 —— 单线程顺利转账")
    r = redis.Redis(connection_pool=POOL)
    r.set("acc:alice", 1000)
    r.set("acc:bob", 100)

    transfer_with_watch("acc:alice", "acc:bob", 200, "T1")

    print(f"\n  Alice 余额 = {r.get('acc:alice')}  (期望 800)")
    print(f"  Bob   余额 = {r.get('acc:bob')}    (期望 300)")


def demo_watch_conflict() -> None:
    section("Demo 4: WATCH 并发冲突 —— 50 线程并发转账,最终一致")
    r = redis.Redis(connection_pool=POOL)
    r.set("acc:alice", 10000)
    r.set("acc:bob", 0)

    THREADS = 50
    AMOUNT_PER = 10

    def worker(i: int) -> None:
        transfer_with_watch("acc:alice", "acc:bob", AMOUNT_PER, f"T{i}")

    threads = [threading.Thread(target=worker, args=(i,)) for i in range(THREADS)]
    t0 = time.time()
    for t in threads: t.start()
    for t in threads: t.join()
    elapsed = time.time() - t0

    alice = int(r.get("acc:alice"))
    bob = int(r.get("acc:bob"))
    expected_total = 10000
    print(f"\n  耗时 {elapsed:.2f}s")
    print(f"  Alice + Bob = {alice + bob}  (期望 {expected_total})")
    print(f"  Bob 收到    = {bob}            (期望 {THREADS * AMOUNT_PER})")
    if alice + bob == expected_total and bob == THREADS * AMOUNT_PER:
        print("  ✅ 余额守恒 —— WATCH 乐观锁工作正常")
    else:
        print("  ❌ 出现不一致")

    r.delete("acc:alice", "acc:bob")


if __name__ == "__main__":
    try:
        demo_basic_transaction()
        demo_exec_runtime_error()
        demo_watch_optimistic_lock()
        demo_watch_conflict()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
python
"""
Ch7 配套代码 2 / 4 —— Pipeline / 单条 / 事务 三种方式压测

演示:
  1. 1 万次 SET,单条调用 —— 每次都走一次 RTT
  2. 1 万次 SET,Pipeline (batch=100) —— 100 倍少 RTT
  3. 1 万次 SET,MULTI/EXEC 整体事务 —— 也省 RTT 但服务端要排队
"""

import time
import redis

HOST, PORT = "127.0.0.1", 6379
N = 10_000          # 总命令数
BATCH = 100         # Pipeline / Tx 批大小


def make_client() -> redis.Redis:
    return redis.Redis(host=HOST, port=PORT, decode_responses=True)


def section(title: str) -> None:
    print("\n" + "=" * 62)
    print(title)
    print("=" * 62)


def warmup() -> None:
    r = make_client()
    r.flushdb()
    for _ in range(50): r.ping()


def bench_single() -> float:
    r = make_client()
    r.flushdb()
    t0 = time.perf_counter()
    for i in range(N):
        r.set(f"bench:s:{i}", "v")
    return time.perf_counter() - t0


def bench_pipeline(transaction: bool = False) -> float:
    r = make_client()
    r.flushdb()
    t0 = time.perf_counter()
    pipe = r.pipeline(transaction=transaction)
    for i in range(N):
        pipe.set(f"bench:p:{i}", "v")
        if (i + 1) % BATCH == 0:
            pipe.execute()
            pipe = r.pipeline(transaction=transaction)
    pipe.execute()
    return time.perf_counter() - t0


def bench_full_transaction() -> float:
    """整个 1 万条全塞一个 MULTI/EXEC 里"""
    r = make_client()
    r.flushdb()
    t0 = time.perf_counter()
    pipe = r.pipeline(transaction=True)
    for i in range(N):
        pipe.set(f"bench:t:{i}", "v")
    pipe.execute()
    return time.perf_counter() - t0


def fmt_row(name: str, sec: float) -> str:
    qps = N / sec if sec > 0 else float("inf")
    return f"  {name:<32} {sec*1000:8.1f} ms     {qps:>10,.0f} qps"


if __name__ == "__main__":
    try:
        warmup()
        section(f"Pipeline 性能压测  N={N}  batch={BATCH}")

        results = []
        results.append(("① 单条 SET(每次一次 RTT)",     bench_single()))
        results.append(("② Pipeline 非事务",              bench_pipeline(False)))
        results.append(("③ Pipeline + 事务(默认)",       bench_pipeline(True)))
        results.append(("④ 单个超大事务(1 万条全打包)",   bench_full_transaction()))

        print()
        print(f"  {'方式':<32} {'耗时':>10}     {'吞吐':>14}")
        print("  " + "-" * 60)
        baseline = results[0][1]
        for name, sec in results:
            print(fmt_row(name, sec))

        print("\n  📊 加速比(相对于「单条」):")
        for name, sec in results:
            speed = baseline / sec if sec > 0 else float("inf")
            print(f"    {name:<32} ×{speed:>5.1f}")

        print("\n  💡 结论:")
        print("    - Pipeline 主要省 RTT,本机 50ms RTT 都能让 1 万次提速 50x")
        print("    - 大事务(④)也省 RTT 但服务端要单线程排队执行 + 占内存")
        print("    - 生产推荐 Pipeline + 中等批大小(100~1000)")

        make_client().flushdb()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
python
"""
Ch7 配套代码 3 / 4 —— Lua 库存扣减实战

演示:
  1. 用纯客户端「GET → 判断 → DECR」会怎么超卖
  2. 用 Lua 原子脚本扣减,1000 并发也精确不超卖
  3. 演示 EVALSHA 缓存、KEYS / ARGV 用法
"""

import threading
import time
import redis

POOL = redis.ConnectionPool(host="127.0.0.1", port=6379, decode_responses=True)


def section(title: str) -> None:
    print("\n" + "=" * 62)
    print(title)
    print("=" * 62)


# ================ Lua 脚本 ================
# KEYS[1]: 库存 key
# ARGV[1]: 要扣的数量
# 返回:    -1 表示库存不足,否则返回剩余库存
LUA_DECR_STOCK = """
local stock = tonumber(redis.call('GET', KEYS[1]) or '0')
local need  = tonumber(ARGV[1])
if stock < need then
    return -1
end
redis.call('DECRBY', KEYS[1], need)
return stock - need
"""

# 滑动窗口限流:1 秒内最多 N 次
# KEYS[1]: 限流 key
# ARGV[1]: 限制次数  ARGV[2]: 窗口秒
LUA_RATE_LIMIT = """
local cur = tonumber(redis.call('GET', KEYS[1]) or '0')
if cur >= tonumber(ARGV[1]) then
    return 0
end
redis.call('INCR', KEYS[1])
if cur == 0 then
    redis.call('EXPIRE', KEYS[1], ARGV[2])
end
return 1
"""


def demo_unsafe_oversell() -> None:
    section("Demo 1: 不用 Lua —— 客户端 GET → 判断 → DECR 会超卖")
    r = redis.Redis(connection_pool=POOL)
    INIT_STOCK = 100
    THREADS = 1000
    BUY_PER = 1

    r.set("stock:item", INIT_STOCK)
    success = []
    lock = threading.Lock()

    def buyer():
        rr = redis.Redis(connection_pool=POOL)
        cur = int(rr.get("stock:item") or 0)
        if cur >= BUY_PER:
            time.sleep(0.0005)  # 模拟一点业务处理时间,制造窗口期
            rr.decrby("stock:item", BUY_PER)
            with lock: success.append(1)

    threads = [threading.Thread(target=buyer) for _ in range(THREADS)]
    for t in threads: t.start()
    for t in threads: t.join()

    final = int(r.get("stock:item"))
    sold = sum(success)
    print(f"  初始库存 {INIT_STOCK}{THREADS} 并发各买 {BUY_PER} 件")
    print(f"  实际成功购买 = {sold}")
    print(f"  剩余库存     = {final}")
    if final < 0 or sold > INIT_STOCK:
        print(f"  ❌ 超卖!(多卖了 {sold - INIT_STOCK} 件)")
    else:
        print(f"  本次没复现超卖(高并发下大概率会发生)")


def demo_lua_safe() -> None:
    section("Demo 2: 用 Lua 原子扣减 —— 1000 并发也不超卖")
    r = redis.Redis(connection_pool=POOL)
    INIT_STOCK = 100
    THREADS = 1000

    r.set("stock:item", INIT_STOCK)
    sha = r.script_load(LUA_DECR_STOCK)
    print(f"  脚本已缓存 SHA = {sha[:16]}...")

    success = []
    lock = threading.Lock()

    def buyer():
        rr = redis.Redis(connection_pool=POOL)
        try:
            res = rr.evalsha(sha, 1, "stock:item", 1)
        except redis.exceptions.NoScriptError:
            res = rr.eval(LUA_DECR_STOCK, 1, "stock:item", 1)
        if res != -1:
            with lock: success.append(1)

    threads = [threading.Thread(target=buyer) for _ in range(THREADS)]
    t0 = time.time()
    for t in threads: t.start()
    for t in threads: t.join()
    elapsed = time.time() - t0

    final = int(r.get("stock:item"))
    sold = sum(success)
    print(f"\n  耗时 {elapsed:.2f}s")
    print(f"  实际成功购买 = {sold}    (期望 {INIT_STOCK})")
    print(f"  剩余库存     = {final}    (期望 0)")
    if sold == INIT_STOCK and final == 0:
        print("  ✅ 完美:原子扣减,零超卖")
    else:
        print("  ❌ 异常")


def demo_rate_limit() -> None:
    section("Demo 3: Lua 滑动计数限流 —— 1 秒内最多 5 次")
    r = redis.Redis(connection_pool=POOL)
    r.delete("rl:user:1001")

    sha = r.script_load(LUA_RATE_LIMIT)
    for i in range(8):
        ok = r.evalsha(sha, 1, "rl:user:1001", 5, 1)
        flag = "✅ 通过" if ok else "🚫 限流"
        print(f"  第 {i+1} 次请求 → {flag}")
        time.sleep(0.05)


def demo_keys_argv() -> None:
    section("Demo 4: KEYS / ARGV 区别 + 多 key 操作")
    r = redis.Redis(connection_pool=POOL)
    script = """
        return {
          'KEYS[1]=' .. KEYS[1],
          'KEYS[2]=' .. KEYS[2],
          'ARGV[1]=' .. ARGV[1],
          'ARGV[2]=' .. ARGV[2],
        }
    """
    res = r.eval(script, 2, "user:1", "user:2", "hello", "world")
    for line in res:
        print(f"  {line}")
    print("\n  💡 numkeys=2 → 前 2 个参数是 KEYS,剩下是 ARGV")
    print("  💡 Cluster 模式必须把 Key 通过 KEYS[] 传,否则路由失败")


if __name__ == "__main__":
    try:
        demo_unsafe_oversell()
        demo_lua_safe()
        demo_rate_limit()
        demo_keys_argv()
        redis.Redis(connection_pool=POOL).delete(
            "stock:item", "rl:user:1001", "user:1", "user:2"
        )
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
python
"""
Ch7 配套代码 4 / 4 —— Stream + Consumer Group 完整示例

演示:
  1. 生产者持续 XADD 订单消息
  2. 同一 Group 内 2 个 worker 协作分摊消费
  3. ACK 机制 + PEL(待确认列表)观察
  4. 模拟 worker-2 崩溃,未 ACK 的消息由 worker-1 通过 XCLAIM 接管
"""

import threading
import time
import random
import redis

POOL = redis.ConnectionPool(host="127.0.0.1", port=6379, decode_responses=True)

STREAM = "orders"
GROUP = "payproc"
N_MESSAGES = 20


def section(title: str) -> None:
    print("\n" + "=" * 62)
    print(title)
    print("=" * 62)


def setup() -> None:
    r = redis.Redis(connection_pool=POOL)
    r.delete(STREAM)
    try:
        r.xgroup_create(STREAM, GROUP, id="$", mkstream=True)
    except redis.ResponseError as e:
        if "BUSYGROUP" not in str(e): raise
    print(f"  ✅ 流 {STREAM} 与消费组 {GROUP} 已就绪")


def producer() -> None:
    r = redis.Redis(connection_pool=POOL)
    for i in range(N_MESSAGES):
        msg_id = r.xadd(STREAM, {
            "order_id": str(1000 + i),
            "user": f"u{random.randint(1,99):02d}",
            "amount": str(random.randint(10, 999)),
        })
        print(f"  [P] XADD {msg_id}  order_id={1000+i}")
        time.sleep(0.1)
    print("  [P] ✅ 生产完毕")


def worker(name: str, fail_rate: float = 0.0, stop_after: int = -1) -> None:
    """
    fail_rate: 模拟「领走但不 ACK」的概率
    stop_after: 处理 N 条后退出(模拟崩溃)。-1 = 不退出
    """
    r = redis.Redis(connection_pool=POOL)
    handled = 0
    while True:
        msgs = r.xreadgroup(
            groupname=GROUP, consumername=name,
            streams={STREAM: ">"}, count=2, block=2000,
        )
        if not msgs:
            print(f"  [{name}] 无新消息 5s,退出")
            return
        for _stream, entries in msgs:
            for msg_id, body in entries:
                handled += 1
                will_fail = random.random() < fail_rate
                tag = "💥 故意不 ACK" if will_fail else "✅ ACK"
                print(f"  [{name}] got {msg_id}  order={body.get('order_id')}  {tag}")
                time.sleep(random.uniform(0.05, 0.2))
                if not will_fail:
                    r.xack(STREAM, GROUP, msg_id)
                if stop_after > 0 and handled >= stop_after:
                    print(f"  [{name}] 🛑 模拟崩溃,已处理 {handled} 条退出")
                    return


def show_pending() -> None:
    section("Demo: XPENDING —— 谁还有没 ACK 的消息?")
    r = redis.Redis(connection_pool=POOL)
    summary = r.xpending(STREAM, GROUP)
    print(f"  pending 总数 = {summary['pending']}")
    if summary["pending"] == 0:
        print("  ✅ 所有消息都已 ACK")
        return
    print(f"  消费者明细 = {summary['consumers']}")

    detail = r.xpending_range(STREAM, GROUP, "-", "+", count=20)
    for d in detail:
        print(f"    {d['message_id']}  consumer={d['consumer']}  idle={d['time_since_delivered']}ms")
    return detail


def demo_claim(detail) -> None:
    """把 worker-2 的死消息转给 worker-1 重新处理"""
    if not detail: return
    section("Demo: XCLAIM —— 把死信转给 worker-1 重新处理")
    r = redis.Redis(connection_pool=POOL)

    dead_ids = [d["message_id"] for d in detail if d["consumer"] != "worker-1"]
    if not dead_ids:
        print("  没有需要 claim 的消息")
        return

    print(f"  XCLAIM {len(dead_ids)} 条 → worker-1(min-idle=0 强制接管)")
    claimed = r.xclaim(STREAM, GROUP, "worker-1", min_idle_time=0, message_ids=dead_ids)
    for msg_id, body in claimed:
        print(f"    worker-1 接管 {msg_id}  order={body.get('order_id')}  → ACK")
        r.xack(STREAM, GROUP, msg_id)


def main() -> None:
    section("Stream Consumer Group 完整流程")
    setup()

    section("Demo: 1 个 Producer + 2 个 Worker 协作消费")
    threads = [
        threading.Thread(target=producer, name="P"),
        threading.Thread(target=worker, args=("worker-1", 0.0, -1), name="W1"),
        threading.Thread(target=worker, args=("worker-2", 0.4, 5), name="W2"),
    ]
    for t in threads: t.start()
    for t in threads: t.join()

    detail = show_pending()
    demo_claim(detail)
    show_pending()

    section("收尾")
    r = redis.Redis(connection_pool=POOL)
    info = r.xinfo_stream(STREAM)
    print(f"  Stream length     = {info['length']}")
    print(f"  Last generated id = {info['last-generated-id']}")
    print(f"  ✅ Demo 完成")
    r.delete(STREAM)


if __name__ == "__main__":
    try:
        main()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")

01_transaction_watch.py ↗ · 02_pipeline_benchmark.py ↗ · 03_lua_stock.py ↗ · 04_stream_consumer_group.py ↗