Celery 太重了?这可能是你一直在找的 asyncio 任务队列

简介: arq是Python原生异步任务队列,基于asyncio与Redis Streams构建,轻量(仅依赖redis-py)、高性能、开箱即用。专为FastAPI等异步框架设计,告别Celery的复杂配置,实现高并发I/O任务的优雅调度。(239字)

你是否在配置 Celery 时,被复杂的 Broker、Backend、Worker 搞得头皮发麻?你是否在 Python async/await 的世界里,发现老牌的任务队列总是显得“格格不入”?

在 Python 异步编程(Asynchronous Programming)日益普及的今天,我们需要一个原生支持 asyncio极度轻量性能强悍的任务队列。

今天,我要向你推荐一个“小而美”的神器——arq

💡 读完本文,你将获得:

  1. 彻底理解 arq 的核心设计理念(小白也能懂)。
  2. 掌握 arq 的安装、配置与实战代码(含优雅写法)。
  3. 深入理解其基于 Redis Streams 的底层原理(高手进阶)。

🧐 What & Why

什么是 arq?

简单来说,arq 是一个基于 Python asyncioRedis 的作业队列(Job Queue)。它的作者是 pydantic 的大神 Samuel Colvin(没错,就是那个写数据验证库的大佬)。

💡 一个生活化的比喻

为了让你秒懂 arq 和 Celery 的区别,我们想象一下“餐厅后厨”的场景:

  • Celery(老牌霸主)
    就像一个五星级酒店的行政总厨。功能极其强大,能做中餐、西餐、日料;有专门的切菜部、炒菜部、传菜部。但是,如果你只想煮一碗泡面,启动这一整套流程就显得非常笨重、繁琐,启动慢,资源消耗大。

  • arq(新晋网红)
    就像一个智能全自动炒菜机器人。它只做一件事:高效地处理订单。它直接连接“菜篮子”(Redis),用最快的速度(asyncio)并行处理几百个订单。它没有复杂的层级,轻便、极速、开箱即用

为什么要选 arq?

  1. 原生异步:天生支持 async/await,与 FastAPI、Sanic 等异步框架是绝配。
  2. 极简依赖:只依赖 redis-py,没有复杂的依赖树。
  3. 高性能:利用 Redis Streams 和 asyncio,并发能力极强。

🛠️ How

1. 快速上手

首先,安装它:

pip install arq

2. 定义 Worker(消费者)

arq 的 worker 定义非常简单,通常我们创建一个 worker.py

# worker.py
import asyncio
from arq import create_pool
from arq.connections import RedisSettings

# 1. 定义具体的任务函数
async def say_hello(ctx, name: str):
    """
    一个简单的异步任务
    :param ctx: 上下文,包含 redis 连接等信息
    :param name: 任务参数
    """
    await asyncio.sleep(1) # 模拟耗时操作
    print(f"👋 Hello {name}, 任务完成!")
    return f"Hello {name}"

# 2. 定义 Worker 配置
class WorkerSettings:
    # Redis 配置
    redis_settings = RedisSettings(host='localhost', port=6379)
    # 注册任务函数
    functions = [say_hello]
    # 并发数
    max_jobs = 10 

    # 💡 优雅写法:生命周期管理
    async def on_startup(self):
        print("🚀 Worker 启动啦!")

    async def on_shutdown(self):
        print("🛑 Worker 停止啦!")

启动 Worker:

arq worker.WorkerSettings

3. 发布任务(生产者)

在你的业务代码(例如 FastAPI 接口)中这样调用:

# main.py
import asyncio
from arq import create_pool
from arq.connections import RedisSettings

async def main():
    # 1. 创建 Redis 连接池
    redis = await create_pool(RedisSettings(host='localhost', port=6379))

    # 2. 入队任务
    # 这里的 'say_hello' 必须和 worker 中注册的函数名一致
    print("📤 正在发送任务...")
    await redis.enqueue_job('say_hello', name='World')

    # 3. 关闭连接
    await redis.close()

if __name__ == '__main__':
    asyncio.run(main())

⚠️ 避坑指南

  1. 不要在任务中执行同步阻塞代码
    arq 是基于 asyncio 的,如果你在任务里写了 time.sleep(10) 或者同步的 requests.get(),整个 Worker 的事件循环会被卡死!
    • ✅ 正确:await asyncio.sleep(10)httpx.AsyncClient
    • ❌ 错误:time.sleep(10)
  2. 参数序列化陷阱
    arq 默认使用 pickle 进行序列化,这很方便,支持几乎所有 Python 对象。但在跨语言或安全性要求高的场景,建议配置为 JSON 序列化。
  3. 任务去重
    arq 支持 _job_id 参数。如果你希望同一个任务 ID 在同一时间只执行一次(防止重复扣款等),可以在入队时指定 _job_id

多聊一聊

对于资深玩家,你可能想问:arq 到底是怎么利用 Redis 实现队列的?

1. Redis Streams 的妙用

在 arq 的早期版本中,它使用 Redis List (LPUSH/BRPOP)。但从 v0.16 开始,它全面拥抱了 Redis Streams(Redis 5.0+ 新特性)。

  • 可靠性:Streams 支持 Consumer Groups(消费者组)。这意味着如果一个 Worker 拿走了任务但崩溃了(没有 ACK),任务不会丢失,会保留在 Pending List 中,等待被其他 Worker 认领(Claim)。
  • 持久化:相比 Pub/Sub 的“发后即焚”,Streams 是持久化的日志结构。

2. Lua 脚本保证原子性

arq 大量使用了 Redis Lua 脚本来保证操作的原子性。
比如“入队”操作,不仅要写入 Stream,可能还需要检查是否有延迟任务(ZSET),或者是否有唯一性约束。arq 将这些逻辑封装在 Lua 脚本中,一次网络往返即可完成,既快又安全。

3. 延迟任务 (Deferred Jobs)

arq 如何实现 enqueue_job('task', _defer_by=60)(延迟60秒执行)?

  • 它并不是让 Worker sleep 60秒。
  • 它是将任务放入 Redis 的 Sorted Set (ZSET) 中,score 是执行时间戳。
  • Worker 内部有一个协程轮询这个 ZSET,时间一到,立马把任务“搬运”到 Streams 队列中供消费。

总结

  • 定位:轻量、异步、基于 Redis。
  • 核心组件:WorkerSettings, enqueue_job, RedisSettings。
  • 底层:Redis Streams (队列) + ZSET (延迟/定时) + Lua (原子性)。
  • 适用场景:高并发 I/O 任务、轻量级微服务、FastAPI 背景任务。

如果你的业务中既有CPU 密集型任务(如视频转码,会阻塞 EventLoop),又有I/O 密集型任务(如发邮件),你会如何设计 arq 的 Worker 架构?是混合部署还是拆分部署?

欢迎在评论区留下你的架构方案!

相关文章
|
7月前
|
人工智能 Rust 安全
OpenClaw Skills深度玩转指南:2026年阿里云部署OpenClaw/Clawdbot+浏览器与邮件技能实战
如果说OpenClaw与大模型的组合是打造智能AI助理的“大脑”,那么Skills就是赋予它行动能力的“双手”。作为OpenClaw生态的核心扩展模块,Skills通过标准化功能封装,让AI助手能够自主完成网页浏览、信息检索、邮件管理等实际操作,彻底打破“只会说不会做”的局限。2026年最新版OpenClaw已默认集成浏览器操作插件agent-browser v0.2.0,同时支持从Clawhub技能库扩展更多实用功能。本文将先介绍阿里云OpenClaw(原Clawdbot)的快速部署步骤,再详细拆解默认Skills的实战场景与新技能安装方法,搭配可直接复用的指令与代码,让新手也能快速解锁AI
3098 1
|
5月前
|
存储 人工智能 BI
Coze开发自能体的费用
Coze(扣子)2026年全面升级计费体系,分国内版(coze.cn,订阅+资源包)与国际版(coze.com,点数制)。国内版含免费/进阶/高阶/企业四档;国际版按模型消耗Credits,GPT-4o等高价、GPT-3.5等低价。另含API调用、商业流量、知识库存储等潜在费用。个人测试选免费版,商用推荐进阶版。(239字)
|
3月前
|
存储 人工智能 缓存
AI 智能体的开发技术
AI智能体是具备感知、思考、规划与执行闭环的真正自主系统,远超简单问答。其开发需融合五大核心技术:智能体编排框架(如LangGraph、CrewAI)、工具调用与安全沙盒、记忆存储(向量库+检查点)、大模型托管平台及可观测性调试工具,实现企业级可靠落地。(239字)
|
2月前
|
NoSQL 安全 Go
Redis 分布式锁的 5个坑,真的是又大又深!!
本文深入剖析Redis分布式锁在Go微服务中的五大演进:从原子性陷阱、安全解锁、看门狗续期,到Go特有重入设计与Fencing Token终极防御。结合Functional Options、context管控、Lua脚本及工程化实践,揭示生产级锁的暗坑与破局之道。(239字)
171 0
|
5月前
|
人工智能
HappyHorse上架阿里云百炼,开启AI视频创作新纪元
阿里巴巴旗下AI视频模型HappyHorse上线阿里云百炼平台,支持文生视频、图生视频、多图参考生成与编辑,具备15秒多镜头叙事、多画幅适配及1080P超清输出能力,大幅降低创作门槛,赋能全场景视频生产。
|
6月前
|
人工智能 IDE API
阿里云百炼Coding Plan 可以接入VS Code、Trae、Cursor等IDE吗?
阿里云百炼Coding Plan支持VS Code、Trae、Cursor三款VS Code内核的AI编辑器,需安装Qwen Code等官方IDE插件接入。仅限AI编程工具或OpenClaw类Agent中使用,不支持API直调或工作流平台。(239字)
|
9月前
|
负载均衡 应用服务中间件 Nacos
Nacos配置中心
本文详细讲解了Nacos作为配置中心的核心功能与实践应用,涵盖配置管理、热更新、共享配置及优先级规则,并通过搭建Nacos集群实现高可用部署,帮助开发者掌握微服务环境下配置的集中化管理方案。
 Nacos配置中心
|
8月前
|
自然语言处理 数据挖掘 测试技术
Qwen3-VL-Embedding系列上新:探索统一多模态表征与排序
2025年6月,Qwen3-VL-Embedding与Qwen3-VL-Reranker开源,基于Qwen3-VL打造,支持文本、图像、视频等多模态检索与跨模态理解,具备统一表示学习、高精度重排序能力,广泛适用于全球化多语言场景,助力高效多模态信息检索。
2847 5
|
C语言
【C语言】全局搜索变量却找不到定义?原来是因为宏!
使用条件编译和 `extern` 来管理全局变量的定义和声明是一种有效的技术,但应谨慎使用。在可能的情况下,应该优先考虑使用局部变量、函数参数和返回值、静态变量或者更高级的封装技术(如结构体和类)来减少全局变量的使用。
471 5