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 返回结果 / 状态 / 异常堆栈 |
✔ 再三重申:
- 即时任务 → XADD 丢进 Stream;
- 延迟任务(
enqueue_in/enqueue_at)→ 先放进 ZSet,Redis 自己没有延迟能力,是 ARQ Worker 轮询 ZSet,到期了 XADD 转发进 Stream。
二、最小可跑示例
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()
四条命令一一对应你手搓过的原生命令:
move_delayed_jobs_from_zset_to_stream—— ZSet 模拟延迟,XADD 进 Stream;XREADGROUP GROUP ... >—— 拉新任务,进 pending;XACK + XDEL—— 完成清理 Stream 消息;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} 保留多久,到期自动过期 |
五、⚠️ 基于原理才懂的坑点
- Stream 消息不会自动清理,但 ARQ 执行完主动 XDEL。 如果你在
redis-cli里看到 Stream key 里堆了大量旧消息,说明 Worker 没正常 XDEL(任务卡住 / 崩溃没跑完)——这是排查信号。 job_timeout不是「任务最大执行时间」,是「pending 闲置超时」。 Worker 卡住超过它,任务就被别的 Worker XCLAIM 抢走重跑——会出现同一个任务并行跑两遍。经典坑。解法:任务内部自己做幂等,或调大job_timeout。- 延迟任务的精度取决于 Worker 轮询频率。 ZSet 是 Worker 循环扫描的,Redis 本身没有定时器,延迟有几十~几百 ms 误差。
- 消费组是 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 友好 |
课程复盘路径
redis-cli实操全部 Stream 命令:XADD XGROUP XREADGROUP XACK XPENDING XCLAIM XDEL XTRIM;redis-py手写极简消费组(第 3 章),跑通生产-消费-ack-故障重跑;- 看懂 ARQ 的 key 分布:Stream 存即时、ZSet 存延迟、Hash 存结果;
- 跑 ARQ 最小 demo,同时开
redis-cli monitor看底层命令流; - 读 ARQ Worker 主循环源码,对应伪代码理解 ZSet 延迟转发 + Stream 消费 + XCLAIM 抢超时。