主题
第 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 给出的理由:
- Redis 命令大部分错误是编程错误(类型用错、参数用错),生产环境不该出现,不该为了少数 bug 让所有事务付出回滚的复杂度代价。
- 没有回滚,Redis 事务的实现可以保持简单且更快。
- 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-py的pipeline(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 注意事项
- 批大小别太大:单批超过几 MB 会撑爆服务端 query buffer,建议 100 ~ 1000 一批。
- 响应必须按顺序读完:客户端拿到响应数组的顺序就是发出命令的顺序。
- Cluster 下要 KEY 同槽:跨槽 Key 无法在同一个 Pipeline 里发(要拆分到不同节点)。
- 不影响其他客户端: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 个参数是 KeyKEYS[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-07.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.py | 1 万次 SET 三种方式对比:单条 / Pipeline / 事务 |
03_lua_stock.py | Lua 库存扣减: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 事务的真实行为是:
- 入队阶段语法错 → 整个事务在 EXEC 时被拒绝(
EXECABORT),算是「全有或全无」。 - 执行阶段类型错(如对字符串 INCR) → 该命令报错,前后命令照常执行,不会回滚。
为什么这样设计? Redis 作者 antirez 的理由:
- 大部分执行错误是程序 bug(用错类型/参数),生产应该测试时就发现,而不是靠回滚兜底。
- 不做回滚,事务实现简单、AOF 持久化无需 undo log、整体更快。
- 「快」是 Redis 第一原则,为少数错误情况付出复杂度不划算。
加分项:提到 ACID 中 Redis 事务保证的是 C(命令一致性)+ I(单线程隔离),A 是「弱原子」,D 取决于 AOF 配置。
Q2:WATCH 是怎么实现乐观锁的?
考察点:CAS 机制理解 + 源码思想。
标准答案:
WATCH 是基于 CAS(Compare-And-Swap)的乐观锁实现。流程:
- 客户端
WATCH key1 key2→ 服务端在db.watched_keys记录{key → 客户端列表}。 - 任何其他客户端的写命令在执行时,会查
watched_keys,把所有 watch 了该 key 的客户端打上CLIENT_DIRTY_CAS标记。 - 客户端
MULTI→命令排队→EXEC。 - 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-py 的 pipe = 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 的硬伤:
- 不持久化:消息发出去就被立即推给在线订阅者,没有 backlog,发完即弃。
- 订阅者掉线丢消息:那段时间发的消息全错过,重连后看不到历史。
- 无 ACK 无重试:消费者收完没处理或处理失败,没机制重投。
- 慢消费者会被踢:服务端 output buffer 满 → 强制断开(
client-output-buffer-limit pubsub)。 - 扇出语义固定:所有订阅者都收到同一条,做不了「多 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 有什么异同?
考察点:对消息队列模型的理解,区分中间件特性。
标准答案:
相同点:
- 消费分摊:同 Group 内每条消息只投递给一个 Consumer,多 worker 协作。
- 多 Group 扇出:不同 Group 之间消息隔离,每个 Group 都能完整消费。
- 位点(offset)持久化:消费进度记录在服务端,重启不丢。
- 支持回放:可以从指定位置重新读。
不同点:
| 维度 | Redis Stream | Kafka |
|---|---|---|
| 存储模型 | 单机内存(可持久化) | 分布式磁盘日志(segment + index) |
| 分区机制 | 无原生分区,靠多 Stream 模拟 | Topic → Partition,原生分区 |
| 消费分摊单位 | 单条消息(worker 抢) | Partition(Partition 绑定 Consumer) |
| 顺序性 | 单 Stream 全局有序,Group 内不保证单 worker 顺序 | Partition 内有序 |
| ACK 粒度 | 逐条 XACK | offset 提交(一批一起) |
| 消息重投 | XCLAIM 把 PEL 里的转给别人 | rebalance + 从上次 offset 重读 |
| 吞吐 | 数十万 qps(单实例) | 百万 qps(集群) |
| 持久化保证 | 取决于 AOF 配置(fsync) | ISR + acks=all |
| Pull vs Push | Pull(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 ↗