KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
01 · Redis Stream 原生命令:消费者组是核心 — keel 龙骨
这一章只做一件事:让你在 redis-cli 里把 Stream 的命令亲手敲一遍。不碰 Python、不碰 ARQ。理由很直接——ARQ 重度依赖 Stream 的消费者组,消费者组没吃透,ARQ 永远是一团谜。
这一章只做一件事:让你在
redis-cli里把 Stream 的命令亲手敲一遍。不碰 Python、不碰 ARQ。理由很直接——ARQ 重度依赖 Stream 的消费者组,消费者组没吃透,ARQ 永远是一团谜。
一、先建立心智模型:Stream 是什么
Stream = Redis 实现的持久化消息日志流。
- 它是一个 key,类型是
stream; - 每一条消息有 Redis 自动生成的 ID:
毫秒时间戳-序号,例如1759012345678-0; - 消息内容是一组 field-value(类似 hash);
- 最关键的差异:和 List 不同,消息读完不会自动删除;和 Pub/Sub 不同,它落盘持久,断线重连能补读。
| 对比 | List | Stream |
|---|---|---|
| 读走即删? | 是(lpop/rpop) | 否,持久留着 |
| 消费确认(ack) | 自己写 | 消费者组原生支持 |
| 多消费者协调 | 几乎无 | 消费者组 + pending 队列 |
| 断线补读 | 做不到(没了) | 从任意 ID 续读 |
二、基础命令:先会写、会读
① XADD:追加一条消息
# key=mystream;* 让 Redis 自动生成 ID;后面是 field value
XADD mystream * name alice age 22
# 返回:1759013012100-0 ← 这条消息的身份证
「自动 ID」是后面所有断线复原机制的地基——客户端只要记住「我读到哪条 ID 了」,重连后从该 ID 之后接着读即可。
② XREAD:普通读取(注意:没有 ack)
# 从 $ 开始,只监听之后新增的消息;BLOCK 0 = 无限阻塞等待
XREAD BLOCK 0 STREAMS mystream $
$ 表示「只盯新增」。⚠️ XREAD 只是简单读消息,没有 ack、没有消费记录,适合发布订阅,不适合任务队列——因为谁读了、读没读完、读一半崩了怎么办,它一概不知。任务队列必须用下面的消费者组。
三、✨ 核心:消费者组(Consumer Group)
这是整门课最重要的概念,ARQ 的任务重试、故障转移全建立在它之上。
消费者组解决三个问题:一组消费者共同消费同一个 Stream;一条消息只分给组内一个消费者;消息分发出去后放进 pending 队列,客户端必须 XACK 才确认完成——没 ack 就一直挂着,消费者崩了可以转给别人重跑。
① 创建消费组 XGROUP CREATE
XGROUP CREATE mystream mygroup 0 MKSTREAM
# 0 = 从头开始消费;MKSTREAM = stream 不存在就自动建
② 组内读取 XREADGROUP
# mygroup 组的 consumerA 消费者;> 只拿「组内没人领过的新消息」
XREADGROUP GROUP mygroup consumerA COUNT 1 BLOCK 0 STREAMS mystream >
> 这个符号是关键:它代表「尚未分配的新消息」。如果你不用 >,而是传一个具体 ID,读到的就是自己这个消费者 pending 里没 ack 的旧消息。
③ 确认完成 XACK(必须!)
业务处理完,必须 ack,消息才从 pending 移除:
XACK mystream mygroup 1759013012100-0
④ 看 pending 队列 XPENDING
XPENDING mystream mygroup
能看到:未 ack 消息数量、最小/最大 ID、各个消费者卡住几条。如果程序崩溃没来得及 XACK,消息就一直挂在这。
⑤ 转移未确认消息 XCLAIM(故障转移!)
# 超过 5000ms 没 ack,把这条消息抢给 consumerB 重新处理
XCLAIM mystream mygroup consumerB 5000 1759013012100-0
👉 这就是任务失败重试的底层来源。 一个 Worker 卡死或崩溃,超时后别的 Worker 用 XCLAIM 把消息抢回来重跑。ARQ 里的 job_timeout 参数,底层就是 XCLAIM 的超时阈值。
⑥ 删消息 / 修剪 XDEL / XTRIM
Stream 消息默认不会自动删,日志会一直膨胀把内存吃爆:
XDEL mystream 1759013012100-0 # 删单条
XTRIM mystream MAXLEN ~ 1000 # 只保留最近 1000 条
⚠️ 这点很重要:原生 Stream 不会丢消息。ARQ 处理完任务后会 XDEL 删掉这条,否则 Redis 内存持续涨。
四、原生 Stream 任务队列的完整流程(记住它)
1. XADD 生产者推送任务消息到 stream key
2. XGROUP CREATE 预先建好消费组
3. XREADGROUP > Worker 阻塞拉取,消息进入 pending
4. 执行业务逻辑
✅ 成功:XACK + XDEL 删除消息;写结果
❌ 崩溃没 ack:消息留在 pending;超时后别人 XCLAIM 抢回重跑
五、🚫 最大的误区
Redis Stream 本身没有「延迟任务」! 它只能立即消费。延迟是上层应用自己实现的——用 ZSet(有序集合)存「将来要执行的任务」,Worker 每次循环主动扫描 ZSet,到期了手动 XADD 丢进 Stream。这是第 6 章讲 ARQ 延迟任务时的核心认知。
↓ 下一步:02 章 · Stream vs Pub/Sub vs List 选型 —— 到底什么时候用哪一个。