Develop

Celery 分布式任务队列实战:从异步任务到定时调度的完整指南

✎ -- 字 🕐 -- 分钟
字号

Celery 分布式任务队列实战封面

Web 应用里总有些操作不能让用户盯着浏览器等:发送一封邮件、生成一份 PDF 报表、批量处理图片、同步第三方 API。如果这些活儿都塞在主请求里,接口分分钟超时,用户体验也会崩。Celery 是 Python 生态里最受欢迎的分布式任务队列之一,这篇文章从安装、任务定义、Worker 调度、定时任务到生产踩坑,带你跑通一条完整链路。

全文代码可直接复制运行,示例基于 Celery 5.4 与 Redis 7,操作系统以 Linux 为主,Windows 用户注意文中标注的特殊参数即可。

一、为什么需要任务队列

先看一下没有任务队列时的典型场景。用户点击"导出月度报表",后端开始查库、聚合、渲染 Excel,最后返回下载链接。整个过程可能持续几十秒,HTTP 连接一直挂着,中间只要网络抖动一下,用户就要重新点。更麻烦的是,这种同步请求会长期占用 Web 服务器的工作线程,高并发时很容易把连接池打满。

任务队列的思路很简单:把耗时的活包装成一个任务,扔进队列里立刻返回"任务已提交",再由后台 Worker 异步消费。主请求在毫秒级结束,真正的重活交给 Worker 慢慢做。除了削峰,任务队列还能做:

  • 解耦:业务系统与具体执行逻辑分离,任务代码可以独立迭代;
  • 削峰填谷:突发流量先排队,Worker 按自身能力处理,避免数据库被打挂;
  • 定时调度:配合 Celery Beat 实现类 Cron 的周期任务,不用依赖系统 crontab;
  • 失败重试:任务失败自动按策略重试,避免手动补偿,也更容易追踪;
  • 分布式扩展:Worker 可以部署在多台机器上,按负载水平扩容。

Python 生态里可选方案不少:Celery、RQ、Huey、Dramatiq。Celery 功能最全、社区最大、插件最多,但配置也稍复杂。如果你只是偶尔发几封邮件,RQ 可能更轻;如果要做企业级调度、路由、监控,Celery 仍是首选。

二、Celery 核心架构

Celery 架构图

上图是 Celery 最常见的部署形态,包含四个核心角色:

  • Producer(生产者):你的 Web 应用,调用 task.delay()task.apply_async() 把任务发到 Broker;
  • Broker(消息中间件):Redis 或 RabbitMQ,负责存储和分发任务消息;
  • Worker(消费者):实际执行任务的进程,可以水平扩展;
  • Backend(结果后端):Redis、数据库等,用于存储任务结果和状态。

Beat 是独立进程,负责按时间表把定时任务投递到 Broker,本身不执行任务。任务执行仍由 Worker 完成。Producer 和 Worker 之间不直接通信,所有协调都通过 Broker 完成,这也是 Celery 可以水平扩展的关键。

三、Broker 选型:Redis 还是 RabbitMQ

Broker 的选择直接影响可靠性、吞吐量和运维成本。下面是两款最常用 Broker 的对比:

维度RedisRabbitMQ
部署难度极简单,一条命令启动稍复杂,需 Erlang 环境
吞吐量非常高,适合短时任务高,支持更复杂路由
消息持久化RDB/AOF,非专业消息持久化原生支持队列和消息持久化
优先级队列支持有限原生支持
路由能力基本Exchange + Routing Key 非常灵活
运维监控Redis 自身监控自带 Management UI
适用场景中小规模、任务可丢失可重跑金融、电商等对可靠性要求高的场景

个人建议:刚起步用 Redis 足够;任务量上来、对可靠性要求变高后,切到 RabbitMQ。切换成本不高,主要改 broker URL 即可。

四、最小可运行示例

4.1 安装依赖

这里用 Redis 做 Broker 和 Backend,安装 redis 客户端与 Celery:

python -m venv venv
source venv/bin/activate  # Windows 用 venv\Scripts\activate
pip install celery[redis] redis

4.2 创建 Celery 应用

新建 celery_app.py

from celery import Celery

app = Celery(
    'demo',
    broker='redis://localhost:6379/0',
    backend='redis://localhost:6379/0',
    include=['tasks']
)

app.conf.update(
    task_serializer='json',
    accept_content=['json'],
    result_serializer='json',
    timezone='Asia/Shanghai',
    enable_utc=True,
)

再新建 tasks.py

from celery_app import app
import time

@app.task(bind=True, max_retries=3, default_retry_delay=10)
def send_email(self, to: str, subject: str, body: str):
    try:
        print(f"Sending email to {to} ...")
        time.sleep(2)  # 模拟网络请求
        return {"to": to, "status": "sent"}
    except Exception as exc:
        raise self.retry(exc=exc, countdown=60)

@app.task
def generate_report(user_id: int, month: str):
    print(f"Generating report for user {user_id}, month {month}")
    time.sleep(5)
    return f"/reports/{user_id}_{month}.pdf"

4.3 启动 Worker

celery -A celery_app worker -l info -P eventlet

Windows 用户建议加 -P eventlet-P gevent,因为 Celery 默认的 prefork 在 Windows 下不支持。Linux/macOS 可以直接跑默认的 prefork,它基于进程,能利用多核 CPU。

4.4 调用任务

再开一个终端:

from tasks import send_email, generate_report

# 异步发送,立即返回 AsyncResult
result = send_email.delay("alice@example.com", "Hi", "Welcome!")
print(result.id)          # 任务 ID
print(result.status)      # PENDING / SUCCESS / FAILURE

# 生成报表
report = generate_report.delay(42, "2026-07")
print(report.get(timeout=10))  # 阻塞等待结果

五、任务调用高级技巧

delay() 是最简单的调用方式,实际生产中更常用 apply_async(),因为它支持更多参数:

send_email.apply_async(
    args=["bob@example.com", "Order", "Your order is shipped"],
    countdown=300,        # 5 分钟后执行
    eta=datetime(2026, 8, 5, 9, 0),  # 指定时间点
    expires=3600,         # 1 小时内有效
    queue='email',        # 指定队列
    priority=5,           # 优先级(需 Broker 支持)
    retry=True,
    retry_policy={
        'max_retries': 5,
        'interval_start': 30,
        'interval_step': 30,
        'interval_max': 300,
    }
)

5.1 签名(Signatures)与链式任务

有时候任务之间有依赖:先下载图片,再压缩,最后上传 CDN。Celery 用 Signature 可以把任务组合起来:

from celery import chain, group, chord
from tasks import download, compress, upload

# 顺序执行:download -> compress -> upload
pipeline = chain(download.s("url"), compress.s(), upload.s())
pipeline.delay()

# 并行执行 10 个下载
group(download.s(f"url_{i}") for i in range(10)).delay()

# 先并行处理,再汇总
chord(
    (download.s(f"url_{i}") for i in range(10)),
    compress.s()
).delay()

5.2 任务序列化与参数安全

默认 Celery 用 JSON 序列化消息。JSON 安全、可读、跨语言,但缺点是不支持 Python 的 datetime、Decimal、自定义对象。如果任务参数里有复杂对象,需要自定义序列化器,或者更简单:只传基础类型和 ID,让 Worker 自己查对象。

另一个常见坑是 pickle。虽然 pickle 能序列化任意 Python 对象,但它有安全风险,不要接收不可信来源的任务。生产环境建议显式限制 accept_content=['json'],关闭 pickle。

六、Worker 启动与监控

生产环境启动 Worker 不建议直接 celery worker,而是用进程管理器。下面是 systemd 单元文件示例:

[Unit]
Description=Celery Worker
After=network.target

[Service]
Type=forking
User=www
Group=www
EnvironmentFile=/etc/default/celeryd
WorkingDirectory=/var/www/myapp
ExecStart=/var/www/myapp/venv/bin/celery multi start worker1 \
    -A celery_app --logfile=/var/log/celery/worker.log \
    --pidfile=/var/run/celery/worker.pid
ExecStop=/var/www/myapp/venv/bin/celery multi stop worker1 \
    --pidfile=/var/run/celery/worker.pid
ExecReload=/var/www/myapp/venv/bin/celery multi restart worker1 \
    -A celery_app --logfile=/var/log/celery/worker.log \
    --pidfile=/var/run/celery/worker.pid

[Install]
WantedBy=multi-user.target

监控方面,Celery 自带的 flower 非常够用:

pip install flower
celery -A celery_app flower --port=5555

Flower 提供 Web UI,可以实时查看任务状态、Worker 负载、失败重试,还能手动触发和撤销任务。配合 Nginx 反向代理和 Basic Auth,可以直接给团队用。

七、定时任务与 Celery Beat

Beat 是 Celery 的调度器。下面配置每分钟执行一次清理日志任务,每天上午九点发送日报:

from celery.schedules import crontab

app.conf.beat_schedule = {
    'cleanup-every-minute': {
        'task': 'tasks.cleanup_logs',
        'schedule': 60.0,
        'args': (),
    },
    'daily-report-at-9am': {
        'task': 'tasks.send_daily_report',
        'schedule': crontab(hour=9, minute=0),
        'args': (),
    },
}
app.conf.timezone = 'Asia/Shanghai'

启动 Beat:

celery -A celery_app beat -l info

注意:Beat 不要部署多份,否则同一任务会被重复调度。生产环境建议用 celery beat -s /var/run/celery/beat-schedule 指定持久化文件,并用 Supervisor 或 systemd 保证单实例。如果要求高可用,可以用 django-celery-beatredbeat 把调度信息存到数据库或 Redis,支持多 Beat 抢锁。

八、路由、队列与优先级

不同任务对延迟和可靠性的要求不一样。邮件可以慢几秒,但支付回调必须立刻处理。Celery 支持把任务路由到不同队列:

app.conf.task_routes = {
    'tasks.send_email': {'queue': 'email'},
    'tasks.process_payment': {'queue': 'critical'},
    'tasks.generate_report': {'queue': 'default'},
}

# 启动 Worker 只消费指定队列
celery -A celery_app worker -Q critical,email -l info

下面是常见队列设计对照表:

队列任务类型Worker 数量优先级策略
critical支付回调、库存扣减多实例、高并发优先保证低延迟
email注册验证、营销邮件2-4 个可容忍分钟级延迟
report报表生成、数据导出1-2 个错峰执行,避免打满 CPU
default普通业务任务视流量而定默认 FIFO

九、重试、Acknowledgement 与幂等性

Celery 默认在任务开始执行时就确认消息(ACK)。如果 Worker 在执行中途崩溃,任务可能丢失。对可靠性要求高的任务,建议开启 task_acks_late=True,让 Worker 在任务完成后才确认:

app.conf.task_acks_late = True
app.conf.worker_prefetch_multiplier = 1

重试虽然好,但一定要保证任务幂等。也就是说,同一条任务执行两次,结果应该和执行一次一样。例如发送邮件前先在数据库里查一下"是否已发送",或者给任务参数加一个唯一 dedupe key。

@app.task(bind=True, max_retries=3)
def charge_order(self, order_id: str):
    order = Order.get(order_id)
    if order.status == "paid":
        return {"id": order_id, "status": "already_paid"}
    try:
        payment_gateway.charge(order)
        order.mark_as_paid()
        return {"id": order_id, "status": "paid"}
    except PaymentError as exc:
        raise self.retry(exc=exc, countdown=30)

十、与 Web 框架集成

Celery 最常见的用法是和 Flask、Django、FastAPI 集成。以 FastAPI 为例,触发任务非常简单:

from fastapi import FastAPI
from tasks import generate_report

app = FastAPI()

@app.post("/reports")
def create_report(user_id: int, month: str):
    task = generate_report.delay(user_id, month)
    return {"task_id": task.id, "status": "submitted"}

@app.get("/reports/{task_id}")
def get_report_status(task_id: str):
    result = generate_report.AsyncResult(task_id)
    return {
        "task_id": task_id,
        "status": result.status,
        "result": result.result if result.ready() else None
    }

关键点:Web 进程只负责投递任务,不要在里面 .get() 阻塞等待结果。如果用户需要轮询状态,单独给一个接口查 Backend。

十一、测试 Celery 任务

单元测试 Celery 任务时,通常不希望真的连 Redis。可以用 task_always_eager=True 让任务同步执行:

import pytest
from celery_app import app

@pytest.fixture(scope='session')
def celery_config():
    return {
        'broker_url': 'memory://',
        'result_backend': 'cache+memory://',
        'task_always_eager': True,
    }

def test_send_email(celery_app):
    from tasks import send_email
    result = send_email.delay("test@example.com", "Hi", "Body")
    assert result.get()["status"] == "sent"

集成测试可以再开一个真实的 Redis,用 pytest fixture 启动临时 Worker,验证端到端流程。

十二、生产配置速查

下面列出 Celery 里最常用的一批配置项,部署时可以直接参考:

配置项推荐值说明
broker_urlredis://host:6379/0消息中间件连接地址
result_backendredis://host:6379/1结果后端,建议与 Broker 分库
task_serializerjson任务序列化方式,生产环境避免 pickle
task_acks_lateTrue任务完成后才确认消息,防丢失
worker_prefetch_multiplier1每个 Worker 每次只取一条任务
worker_max_tasks_per_child1000子进程处理多少任务后重启,防内存泄漏
result_expires3600结果保留 1 小时后自动清理
task_time_limit3600单个任务硬超时时间(秒)
task_soft_time_limit3000软超时,任务可捕获异常做收尾
timezoneAsia/ShanghaiBeat 调度时区

十三、常见陷阱与最佳实践

问题原因解决方案
任务丢失Worker 崩溃前未 ACK 或中间件配置不当task_acks_late=True;RabbitMQ 用持久化队列
任务重复执行超时重试或 Beat 多实例任务幂等;Beat 只跑单实例
内存泄漏Worker 长时间运行积累对象设置 worker_max_tasks_per_child=1000
结果后端爆库默认 backend 永久保存结果设置 result_expires=3600 或不用 backend
任务参数过大把图片、大 JSON 直接塞进任务参数传对象 ID 或文件路径,让 Worker 自己读取
时区混乱未统一 timezoneenable_utc=True,显式设置 timezone
启动失败Windows 默认 prefork 不支持-P eventlet-P gevent
循环导入Celery 应用与任务互相 import把 app 和 task 分开放,用 include 注册

再补充几条生产经验:

  • 不要用数据库做 Broker:SQLAlchemy 结果后端可以,但 Broker 请用 Redis 或 RabbitMQ,否则性能堪忧。
  • 任务不要返回大对象:结果后端一般存 Redis,返回大对象会占内存。
  • 日志要结构化:结合任务 ID 打日志,排查问题时能串起来。
  • Worker 要可观测:Flower + Prometheus exporter,监控队列深度、失败率、执行耗时。
  • 部署升级要优雅:先停 Beat,等 Worker 消费完再重启,避免任务中断。
  • 参数校验前置:任务函数入口做类型和范围校验,避免无效任务进队列。

十四、总结

Celery 不是银弹,但在 Python 后端里处理异步任务、定时调度、削峰填谷,它仍然是最稳妥的选择。核心要记住三点:第一,任务参数尽量小,大的数据让 Worker 自己拉;第二,任务必须幂等,重试才不会导致业务错乱;第三,根据业务重要性拆分队列,不要把邮件和支付回调塞到同一个 Worker 里。

如果你在本地已经跑通了上面的示例,下一步可以试试把它接到 Flask 或 FastAPI 里,用接口触发 Celery,再用 Flower 观察任务流转。把异步链路跑通之后,很多曾经让你头疼的慢接口问题,都会变得清爽很多。