Flask 集成 Celery:后台任务队列的官方配置模式与实战示例
当你需要处理耗时的上传数据处理、发送邮件等长运行任务时,直接在请求中等待会阻塞整个响应。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.rst 及 examples/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
这段代码做了四件事,逐一对应源码行为:
FlaskTask(Task)子类:重写了__call__,在执行任务函数前自动压入一个 Flask 应用上下文(with app.app_context())。这样任务内部就能使用依赖应用上下文的服务,比如通过current_app获取的配置、数据库连接等。没有这一步,在 worker 进程中运行的任务将无法访问 Flask 的配置和依赖。Celery(app.name, task_cls=FlaskTask):创建 Celery app 对象,命名与 Flask app 同名(便于日志排查),并指定task_cls=FlaskTask使所有任务都继承上下文包裹行为。celery_app.config_from_object(app.config["CELERY"]):Celery 配置从 Flask 配置的CELERY键读取。所有配置项都是普通的 Celery 配置(broker_url、result_backend、task_ignore_result等),命名遵循 Celery 的约定。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 模块并从中查找名为 celery 或 celery_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/src/task_app/init.py:应用工厂 +
celery_init_app,配置了CELERY字典(Redis broker、结果后端、task_ignore_result=True)并注册视图蓝图; - examples/celery/src/task_app/tasks.py:三个任务——
add(相加)、block(睡 5 秒)、process(每秒上报一次进度); - examples/celery/src/task_app/views.py:蓝图路由,POST 提交任务返回
result_id,GET 轮询结果; - examples/celery/src/task_app/templates/index.html:前端页面,用 JavaScript 提交任务并每 500ms 轮询一次结果;
- examples/celery/make_celery.py:供
celery -A make_celery使用的 app 查找入口。
按 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 集成模式可以归纳为四个要点:
- 用
celery_init_app(app)工厂式函数把 Celery 挂到 Flask 上:配置取自app.config["CELERY"],自定义Task子类保证任务在应用上下文中运行,结果存入app.extensions["celery"](src/flask/sansio/app.py 中定义的扩展状态字典); - 任务一律用
@shared_task定义,配合celery_app.set_default(),避免与具体 app 实例绑定、兼容工厂模式和测试; - 请求内用
task.delay(...)提交任务并返回result.id,绝不在路由内阻塞等待; - 用
AsyncResult(id)路由轮询ready/successful/value;向任务只传最小可序列化数据(如 id),在任务内部凭应用上下文重新构建复杂对象。
完整的可运行代码见 examples/celery,文档原文见 docs/patterns/celery.rst。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0622
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00