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 的重试,做个实验:

  1. 在消费者循环里,处理某条消息时 time.sleep(1000)(模拟卡死)或 os._exit(1)(模拟崩溃);
  2. 另开终端看 XPENDING demo_stream demo_group——这条消息挂在 pending,没人 ack;
  3. 用 XCLAIM demo_stream demo_group worker2 5000 <msg_id> 把它抢给另一个消费者重跑。

你会亲眼看到「未确认任务被转移重跑」——ARQ 的 job_timeout 干的就是这事。


↓ 下一步:先插入两个不是队列的 Redis 运行时问题——04 章 · 分布式锁 和 05 章 · 业务锁封装,然后再回来看 ARQ。

进入 keel 阅读