首页
/ Flask 集成 Celery:后台任务队列的官方配置模式与实战示例

Flask 集成 Celery:后台任务队列的官方配置模式与实战示例

2026-09-04 20:58:46作者:董斯意

当你需要处理耗时的上传数据处理、发送邮件等长运行任务时,直接在请求中等待会阻塞整个响应。Flask 官方文档(docs/patterns/celery.rst)给出了一套标准的集成模式:用 Flask 配置驱动 Celery、让任务在应用上下文中运行、通过 delay() 提交任务并轮询结果。本文完整继承该文档的核心内容,并结合仓库中的官方示例(examples/celery)与 Flask 源码实现,逐段讲解每个配置项的作用与底层原理。

为什么需要任务队列

如果应用中有长运行任务——比如处理上传的数据或发送邮件——你不希望在请求期间等待它完成。正确的做法是使用任务队列,把必要的数据发送到另一个进程,让任务在后台运行,而请求立即返回。

Celery 是一个强大的任务队列,既可用于简单的后台任务,也可用于复杂的多阶段程序和定时调度。Flask 文档的这篇模式页聚焦于「如何用 Flask 的方式配置 Celery」,而 Celery 本身的用法需要参考其官方文档(First Steps with Celery)。

Flask 仓库附带了一个基于本文内容的完整示例 examples/celery,演示如何用 JavaScript 提交任务并轮询进度和结果,本文最后会结合该示例做实战拆解。

安装

从 PyPI 安装 Celery,例如使用 pip:

$ pip install celery

官方示例的依赖声明为 celery[redis](见 examples/celery/pyproject.toml),因为示例使用 Redis 作为 broker。若使用其他消息代理(如 RabbitMQ),需安装对应的 kombu 传输依赖。

集成 Celery 与 Flask

你也可以完全不依赖 Flask 地使用 Celery,但通过 Flask 配置来配置它、并让任务访问 Flask 应用会带来很大的便利。Celery 与 Flask 思路相近:都有一个带配置、可注册任务的 app 对象。在创建 Flask app 时,用下面的代码同时创建并配置一个 Celery app(与 docs/patterns/celery.rstexamples/celery/src/task_app/init.py 中的 celery_init_app 完全一致):

from celery import Celery, Task

def celery_init_app(app: Flask) -> Celery:
    class FlaskTask(Task):
        def __call__(self, *args: object, **kwargs: object) -> object:
            with app.app_context():
                return self.run(*args, **kwargs)

    celery_app = Celery(app.name, task_cls=FlaskTask)
    celery_app.config_from_object(app.config["CELERY"])
    celery_app.set_default()
    app.extensions["celery"] = celery_app
    return celery_app

这段代码做了四件事,逐一对应源码行为:

  1. FlaskTask(Task) 子类:重写了 __call__,在执行任务函数前自动压入一个 Flask 应用上下文(with app.app_context())。这样任务内部就能使用依赖应用上下文的服务,比如通过 current_app 获取的配置、数据库连接等。没有这一步,在 worker 进程中运行的任务将无法访问 Flask 的配置和依赖。
  2. Celery(app.name, task_cls=FlaskTask):创建 Celery app 对象,命名与 Flask app 同名(便于日志排查),并指定 task_cls=FlaskTask 使所有任务都继承上下文包裹行为。
  3. celery_app.config_from_object(app.config["CELERY"]):Celery 配置从 Flask 配置的 CELERY 键读取。所有配置项都是普通的 Celery 配置(broker_urlresult_backendtask_ignore_result 等),命名遵循 Celery 的约定。
  4. celery_app.set_default() + app.extensions["celery"] = celery_app:把 Celery app 设为默认 app,使 @shared_task 在每次请求期间都能看到它;同时存入 app.extensions 字典。从 Flask 源码看(src/flask/sansio/app.py),app.extensions 是「a place where extensions can store application specific state」,即专为扩展存放应用级状态而设计的字典,这也是 Flask 官方约定的扩展挂载点。

基础示例:用 Redis 作为通信后端

下面是一个基础的 example.py,配置 Celery 使用 Redis 通信。这里启用了结果后端(result backend),但默认忽略结果——这样只在关心结果的任务上才存储结果:

from flask import Flask

app = Flask(__name__)
app.config.from_mapping(
    CELERY=dict(
        broker_url="redis://localhost",
        result_backend="redis://localhost",
        task_ignore_result=True,
    ),
)
celery_app = celery_init_app(app)

三个配置项的作用:

  • broker_url:消息代理地址,任务消息经由此处从提交方投递到 worker;
  • result_backend:结果后端,任务执行结果存储于此,供之后按 id 查询;
  • task_ignore_result=True:全局默认不保存任务结果(省内存),需要结果的任务单独用 ignore_result=False 覆盖。

启动 worker 与 beat

celery worker 命令指向该模块,它就能找到 celery_app 对象:

$ celery -A example worker --loglevel INFO

还可以运行 celery beat 命令按调度执行任务,调度定义方式参见 Celery 官方文档:

$ celery -A example beat --loglevel INFO

应用工厂模式(Application Factory)

使用 Flask 应用工厂模式时,在工厂内部调用 celery_init_app 即可。它把 Celery app 对象设置到 app.extensions["celery"],之后可以从工厂返回的 Flask app 上取回它:

def create_app() -> Flask:
    app = Flask(__name__)
    app.config.from_mapping(
        CELERY=dict(
            broker_url="redis://localhost",
            result_backend="redis://localhost",
            task_ignore_result=True,
        ),
    )
    app.config.from_prefixed_env()
    celery_init_app(app)
    return app

注意 app.config.from_prefixed_env()(定义于 src/flask/config.py):它从环境变量中读取带应用名(大写、下划线转双下划线)前缀的变量覆盖配置,例如 EXAMPLE__CELERY__BROKER_URL。生产部署时用环境变量注入 broker 地址比写死在代码中更灵活,官方示例的工厂函数(examples/celery/src/task_app/init.py)也调用了它。

make_celery.py:让 celery 命令找到 app

要使用 celery 命令,Celery 需要一个 app 对象,但工厂模式下 app 不再直接存在于模块顶层。解决办法是创建一个 make_celery.py 文件,调用 Flask 应用工厂并从返回的 Flask app 中取出 Celery app:

from example import create_app

flask_app = create_app()
celery_app = flask_app.extensions["celery"]

官方示例中的 examples/celery/make_celery.py 正是这三行(从 task_app 导入 create_app)。然后把 celery 命令指向这个文件:

$ celery -A make_celery worker --loglevel INFO
$ celery -A make_celery beat --loglevel INFO

-A make_celery 的含义是:导入 make_celery 模块并从中查找名为 celerycelery_app 的对象。

定义任务:使用 @shared_task

使用 @celery_app.task 装饰器会直接引用 celery_app 对象,这在工厂模式下不可用(装饰发生在模块导入期,此时 celery_app 变量尚未存在)。此外,被这样装饰的任务会与特定的 Flask/Celery app 实例绑定,如果测试中改变了配置,这种绑定会引发问题。

因此应使用 Celery 的 @shared_task 装饰器。它创建的任务对象会访问当前的「current app」——与 Flask 的蓝图和应用上下文是类似的概念。这正是前文调用 celery_app.set_default() 的原因:set_default() 让 shared task 在运行时能找到默认 Celery app。

一个相加两个数并返回结果的示例任务:

from celery import shared_task

@shared_task(ignore_result=False)
def add_together(a: int, b: int) -> int:
    return a + b

此前我们配置 Celery 默认忽略任务结果。由于想知道这个任务的返回值,显式设置 ignore_result=False。反之,不需要结果的任务(比如发送邮件)就不设置该参数。

官方示例中还有两个更完整的任务(examples/celery/src/task_app/tasks.py),演示了阻塞任务与进度上报:

@shared_task()
def block() -> None:
    time.sleep(5)

@shared_task(bind=True, ignore_result=False)
def process(self: Task, total: int) -> object:
    for i in range(total):
        self.update_state(state="PROGRESS", meta={"current": i + 1, "total": total})
        time.sleep(1)

    return {"current": total, "total": total}

bind=True 会把任务对象自身作为第一个参数 self 传入,从而可以用 self.update_state() 更新任务状态和自定义元数据——这是实现进度条的关键。

调用任务

装饰后的函数成为带后台调用方法的任务对象,最简单的方式是 delay(*args, **kwargs) 方法(更多调用方式参见 Celery 文档)。

运行任务需要一个 Celery worker 在跑(前文已演示如何启动 worker)。

from flask import request

@app.post("/add")
def start_add() -> dict[str, object]:
    a = request.form.get("a", type=int)
    b = request.form.get("b", type=int)
    result = add_together.delay(a, b)
    return {"result_id": result.id}

路由不会立即拿到任务结果——那样做会阻塞响应,违背使用任务队列的初衷。取而代之的是返回运行中任务的 result id,之后用它来获取结果。

获取结果

为获取上面启动的任务的结果,添加另一个路由,接收之前返回的 result id。返回任务是否已完成(ready)、是否成功完成、以及完成时的返回值(或错误):

from celery.result import AsyncResult

@app.get("/result/<id>")
def task_result(id: str) -> dict[str, object]:
    result = AsyncResult(id)
    return {
        "ready": result.ready(),
        "successful": result.successful(),
        "value": result.result if result.ready() else None,
    }

这样就能用第一个路由启动任务,用第二个路由轮询结果。Flask 的请求 worker 不再被阻塞等待任务完成。

AsyncResult(id) 按 id 从结果后端反查任务状态;官方示例的路由实现(examples/celery/src/task_app/views.py)略有细化——在任务就绪后用 result.get() 获取值(失败时会抛出异常,配合 successful() 判断可区分成功值与错误):

@bp.get("/result/<id>")
def result(id: str) -> dict[str, object]:
    result = AsyncResult(id)
    ready = result.ready()
    return {
        "ready": ready,
        "successful": result.successful() if ready else None,
        "value": result.get() if ready else result.result,
    }

向任务传递数据

上面的 add 任务接收两个整数参数。向任务传参时,Celery 必须把参数序列化成可传给其他进程的格式,因此不推荐传递复杂对象。例如 SQLAlchemy 模型对象通常不可序列化,而且它绑定在查询它的 session 上,跨进程传递后也毫无意义。

原则是:只传递在任务内部获取或重建复杂数据所需的最小数据。考虑一个任务:当登录用户请求自己的数据归档时运行。Flask 请求知道当前登录用户,并持有从数据库查出的用户对象——它是按给定 id 查库得到的,所以任务可以照做。传用户的 id,而不是用户对象:

@shared_task
def generate_user_archive(user_id: str) -> None:
    user = db.session.get(User, user_id)
    ...

generate_user_archive.delay(current_user.id)

这也解释了前文 FlaskTask 的意义:任务进程里没有请求上下文,但通过 app.app_context() 可以获得 current_app、配置等应用级资源,从而完成「重新查库、重新构建」这一步。

官方示例实战:examples/celery

仓库中的 examples/celery 把上述模式组装成一个可运行的项目,目录结构为:

examples/celery/README.md 的运行步骤(依赖见 examples/celery/requirements.txt,其中锁定了 celery[redis]==5.2.7):

$ python3 -m venv .venv
$ . ./.venv/bin/activate
$ pip install -r requirements.txt && pip install -e .
$ celery -A make_celery worker --loglevel INFO

在另一个终端激活同一虚拟环境,运行 Flask 开发服务器:

$ . ./.venv/bin/activate
$ flask -A task_app run --debug

打开 http://localhost:5000/ 即可用表单提交任务。页面上的 JavaScript 逻辑(index.html 中)展示了三种典型交互形态:

  • Add:提交后持续轮询,直到 ready 为真,显示最终结果值;
  • Block:任务本身要 5 秒,但 HTTP 响应立即返回,直观演示「请求不等待任务」这一核心卖点;
  • Process:轮询过程中 ready 为假时,value 携带 update_state 写入的进度元数据 {current, total},页面据此显示 3 / 10 这样的进度,任务完成后显示完成。

轮询核心逻辑是一个自递归的 setTimeout(poll, 500):每次 fetch /tasks/result/<id>,若 ready 为假则 500ms 后再轮询,若成功则结束,失败则输出到 console。

小结

Flask 官方推荐的 Celery 集成模式可以归纳为四个要点:

  1. celery_init_app(app) 工厂式函数把 Celery 挂到 Flask 上:配置取自 app.config["CELERY"],自定义 Task 子类保证任务在应用上下文中运行,结果存入 app.extensions["celery"]src/flask/sansio/app.py 中定义的扩展状态字典);
  2. 任务一律用 @shared_task 定义,配合 celery_app.set_default(),避免与具体 app 实例绑定、兼容工厂模式和测试;
  3. 请求内用 task.delay(...) 提交任务并返回 result.id,绝不在路由内阻塞等待;
  4. AsyncResult(id) 路由轮询 ready / successful / value;向任务只传最小可序列化数据(如 id),在任务内部凭应用上下文重新构建复杂对象。

完整的可运行代码见 examples/celery,文档原文见 docs/patterns/celery.rst

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
527
590
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
904
1.82 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
889
5.78 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.52 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.33 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
981
502
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384