KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
03 · 用 redis-py 手写最小消费组(先别碰 ARQ) — keel 龙骨
第 1 章你在 redis-cli 里把命令敲熟了。这一章换 Python:用 redis-py 手写一个最小消费者组——生产者 + 消费者循环 + ack。做到能跑通「生产 → 消费 → ack → 故障重跑」全流程。这一步是理解 ARQ 的脚手架:等你手搓过一遍,第 6 章看 ARQ 源码会像看自己写的代码。
第 1 章你在
redis-cli里把命令敲熟了。这一章换 Python:用redis-py手写一个最小消费者组——生产者 + 消费者循环 + ack。做到能跑通「生产 → 消费 → ack → 故障重跑」全流程。这一步是理解 ARQ 的脚手架:等你手搓过一遍,第 6 章看 ARQ 源码会像看自己写的代码。
一、前置:装库、连 Redis
pip install redis
import redis
r = redis.Redis(host="localhost", port=6379, decode_responses=True)
STREAM_KEY = "demo_stream"
GROUP = "demo_group"
CONSUMER = "worker1"
decode_responses=True 让返回的是字符串而不是 bytes,演示阶段省心。
二、生产者:发一条任务
# 组已存在会报错,捕获一下静默跳过
try:
r.xgroup_create(STREAM_KEY, GROUP, id="0", mkstream=True)
except redis.exceptions.ResponseError:
pass
msg_id = r.xadd(STREAM_KEY, {"task": "generate_report", "payload": "user-42"})
print(f"send task id: {msg_id}")
这就对应第 1 章的 XADD + XGROUP CREATE。
三、消费者:阻塞读 + ack(核心循环)
while True:
# 阻塞读取组内「新消息」(符号 >)
res = r.xreadgroup(
groupname=GROUP,
consumername=CONSUMER,
count=1,
block=0, # 0 = 无限阻塞,没有消息就等
streams={STREAM_KEY: ">"},
)
if not res:
continue
stream_data = res[0][1] # [(stream_key, [(id, fields), ...])]
for msg_id, data in stream_data:
print(f"receive {msg_id}, data={data}")
# —— 这里写你的业务:生成报告、发邮件、调模型……
# 模拟:time.sleep(2)
# ✅ 业务完成,ack 让消息离开 pending
r.xack(STREAM_KEY, GROUP, msg_id)
# ✅ 删除消息,避免 Redis 内存持续涨
r.xdel(STREAM_KEY, msg_id)
把这段跑起来,再开另一个终端用 XADD 发消息,你会看到 worker 实时消费、ack、删消息。这就是一个极简版任务队列的核心。
四、🧠 重点:你刚手搓了 ARQ 的心脏
对比第 1 章的原生命令,ARQ 只是在上面这层做了封装:
| 你手写的 | ARQ 帮你封装的 |
|---|---|
xadd(... {"task": ...}) |
job 序列化(pickle)后 XADD |
block=0 阻塞读 > |
Worker 主循环的阻塞拉取 |
| 没写的超时重跑 | job_timeout → 底层 XCLAIM 抢超时卡住的任务 |
| 没写的延迟 | ZSet 存延迟任务,到期 XADD 进 Stream |
| 没写的结果保存 | 结果写进 Hash(arq:result:{job_id}) |
| 没写的重试次数 | 重试计数 + 超过上限标记失败 |
换句话说:你刚才写的 20 行,就是 ARQ 最朴素的内核。 它没魔法,只是把「超时重跑、延迟、结果、重试上限、Worker 管理」这些工程细节补齐全了。
五、故意制造一次故障,验证 XCLAIM
想真的理解 ARQ 的重试,做个实验:
- 在消费者循环里,处理某条消息时
time.sleep(1000)(模拟卡死)或os._exit(1)(模拟崩溃); - 另开终端看
XPENDING demo_stream demo_group——这条消息挂在 pending,没人 ack; - 用
XCLAIM demo_stream demo_group worker2 5000 <msg_id>把它抢给另一个消费者重跑。
你会亲眼看到「未确认任务被转移重跑」——ARQ 的 job_timeout 干的就是这事。
↓ 下一步:先插入两个不是队列的 Redis 运行时问题——04 章 · 分布式锁 和 05 章 · 业务锁封装,然后再回来看 ARQ。