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 实现的持久化消息日志流。

对比 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 选型 —— 到底什么时候用哪一个。

进入 keel 阅读