KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA

06 · ARQ 源码与使用:看懂它在 Redis 里写了什么 — keel 龙骨

第 3 章你手搓过最小消费组。这一章回到主线:ARQ(基于 Redis 的 Python 异步任务队列)到底在 Redis 里写了什么?它的 Worker 主循环是怎么跑的?重点不是会调 API,而是打开 redis-cli monitor,亲眼看到它发出的每一条 Redis 命令——那一刻 ARQ 就不再是黑盒。

第 3 章你手搓过最小消费组。这一章回到主线:ARQ(基于 Redis 的 Python 异步任务队列)到底在 Redis 里写了什么?它的 Worker 主循环是怎么跑的?重点不是会调 API,而是打开 redis-cli monitor,亲眼看到它发出的每一条 Redis 命令——那一刻 ARQ 就不再是黑盒。


一、先纠偏:ARQ 底层就是 Redis Stream + ZSet

ARQ 全部任务、延迟任务、结果、job 状态都活在 Redis 里,没有独立数据库。关键约定(queue_name 默认 arq):

Redis 结构 key 放什么
Stream arq:stream:{queue_name} 即时任务队列(消费组在此)
消费组 arq:group:{queue_name} Stream 对应的消费者组
ZSet arq:z:{queue_name} 延迟任务:score = 计划执行时间戳
Hash arq:result:{job_id} job 返回结果 / 状态 / 异常堆栈

✔ 再三重申:

二、最小可跑示例

Worker 端 worker.py

import asyncio
from arq import create_pool, Worker, cron
from arq.connections import RedisSettings

async def generate_report(ctx, user_id: str):
    print(f"running report for {user_id}")
    return f"done:{user_id}"

async def startup(ctx): ...
async def shutdown(ctx): ...

class WorkerSettings:
    redis_settings = RedisSettings(host="localhost")
    functions = [generate_report]      # 注册可执行的任务函数
    on_startup = startup
    on_shutdown = shutdown
    job_timeout = 10                    # 对应底层 XCLAIM 的超时阈值!

启动:arq worker.WorkerSettings

客户端 client.py

async def main():
    redis = await create_pool(RedisSettings(host="localhost"))
    job = await redis.enqueue_job("generate_report", "user-42")
    print(f"job id {job.job_id}")
    result = await job.result()         # 阻塞等结果
    print("result:", result)
    await redis.close()

asyncio.run(main())

三、🧠 Worker 主循环(把底层打通)

打开 ARQ 源码 arq/worker.py 的 run(),伪代码对应第 1、3 章的原生命令:

while worker_running:
    # step1:扫描 ZSet 里的延迟任务,到期的 XADD 推入 Stream
    move_delayed_jobs_from_zset_to_stream()

    # step2:XREADGROUP 阻塞从 Stream 拉新任务(符号 >)
    jobs = xreadgroup(..., ">")
    if jobs:
        run_task()                      # 进入 pending,执行 async 任务
        # 成功:XACK + XDEL 删除 Stream 消息;写结果 Hash
        # 异常:按重试设置,重投或标记失败写结果 Hash

    # step3:扫描 pending,XCLAIM 超时卡住的任务(job_timeout)
    claim_stale_pending_jobs()

四条命令一一对应你手搓过的原生命令:

  1. move_delayed_jobs_from_zset_to_stream —— ZSet 模拟延迟,XADD 进 Stream;
  2. XREADGROUP GROUP ... > —— 拉新任务,进 pending;
  3. XACK + XDEL —— 完成清理 Stream 消息;
  4. XCLAIM —— 抢超时卡住任务,就是 job_timeout 的底层。

四、ARQ 参数 ↔ Redis Stream 底层映射表

ARQ 参数 Redis 底层
job_timeout XCLAIM 的超时毫秒阈值;Worker 卡死超过它就任务被抢走重跑
queue_name Stream key 后缀 arq:stream:{queue_name}
每个 Worker 实例 消费组里的消费者名(启动时随机生成)
enqueue_in / enqueue_at 写入 ZSet arq:z:{queue_name},Worker 轮询到期转发进 Stream
keep_result 结果 Hash arq:result:{job_id} 保留多久,到期自动过期

五、⚠️ 基于原理才懂的坑点

  1. Stream 消息不会自动清理,但 ARQ 执行完主动 XDEL。 如果你在 redis-cli 里看到 Stream key 里堆了大量旧消息,说明 Worker 没正常 XDEL(任务卡住 / 崩溃没跑完)——这是排查信号。
  2. job_timeout 不是「任务最大执行时间」,是「pending 闲置超时」。 Worker 卡住超过它,任务就被别的 Worker XCLAIM 抢走重跑——会出现同一个任务并行跑两遍。经典坑。解法:任务内部自己做幂等,或调大 job_timeout。
  3. 延迟任务的精度取决于 Worker 轮询频率。 ZSet 是 Worker 循环扫描的,Redis 本身没有定时器,延迟有几十~几百 ms 误差。
  4. 消费组是 Worker 启动时自动 XGROUP CREATE,不用你手动建。

六、📌 强烈建议:开 monitor 亲眼验证

新开终端,提交任务、Worker 消费时,你会亲眼看到 ARQ 发出的 XADD / XREADGROUP / XACK / XCLAIM——底层全貌一次看清:

redis-cli monitor

七、横向对比收尾

队列 底层 可靠性
Celery(Redis broker) Redis List 无原生 ack,可靠性弱
RQ Redis List 无消费组 ack
ARQ Stream 消费组 + ZSet 延迟 ack + XCLAIM 故障转移,asyncio 友好

课程复盘路径

  1. redis-cli 实操全部 Stream 命令:XADD XGROUP XREADGROUP XACK XPENDING XCLAIM XDEL XTRIM;
  2. redis-py 手写极简消费组(第 3 章),跑通生产-消费-ack-故障重跑;
  3. 看懂 ARQ 的 key 分布:Stream 存即时、ZSet 存延迟、Hash 存结果;
  4. 跑 ARQ 最小 demo,同时开 redis-cli monitor 看底层命令流;
  5. 读 ARQ Worker 主循环源码,对应伪代码理解 ZSet 延迟转发 + Stream 消费 + XCLAIM 抢超时。

进入 keel 阅读