在Flask应用中处理发信、大文件生成、数据批量计算等耗时操作时,同步执行会导致请求长时间等待,用户需要等待任务完成才能收到响应,严重影响使用体验。使用Celery可以将这些任务放到后台异步执行,通过消息队列分发任务,还能随时查询任务的执行状态和结果。

环境准备与依赖安装
首先需要安装必要的依赖包,这里选择Redis作为消息队列和结果存储后端,你也可以根据需求替换为RabbitMQ等其他消息队列。
pip install flask celery redis
基础配置与Celery实例初始化
在Flask项目中需要先初始化Flask应用,再配置Celery的相关参数,让Celery能够和Flask应用协同工作。
from flask import Flask
from celery import Celery
# 初始化Flask应用
app = Flask(__name__)
# Flask配置
app.config['CELERY_BROKER_URL'] = 'redis://127.0.0.1:6379/0' # 消息队列地址
app.config['CELERY_RESULT_BACKEND'] = 'redis://127.0.0.1:6379/0' # 结果存储地址
# 初始化Celery实例
def make_celery(app):
celery = Celery(
app.import_name,
broker=app.config['CELERY_BROKER_URL'],
backend=app.config['CELERY_RESULT_BACKEND']
)
# 让Celery任务能访问Flask应用的上下文
celery.conf.update(app.config)
return celery
celery = make_celery(app)
定义异步任务
使用@celery.task装饰器定义异步任务,这里以耗时发信任务为例,模拟发信的耗时操作。
import time
@celery.task(bind=True)
def send_email_task(self, receiver, subject, content):
"""异步发信任务"""
try:
# 模拟发信耗时,实际场景中替换为真实的发信逻辑
self.update_state(state='PROGRESS', meta={'progress': 30})
time.sleep(2)
self.update_state(state='PROGRESS', meta={'progress': 70})
time.sleep(2)
# 假设发信成功
return {'status': 'success', 'receiver': receiver, 'msg': '发信成功'}
except Exception as e:
# 任务执行失败返回错误信息
return {'status': 'failed', 'error': str(e)}
提交异步任务到消息队列
在Flask的接口中调用异步任务,使用delay或者apply_async方法提交任务到消息队列,接口会立即返回任务ID,不会等待任务执行完成。
from flask import jsonify, request
@app.route('/send_email', methods=['POST'])
def trigger_send_email():
"""触发异步发信的接口"""
data = request.get_json()
receiver = data.get('receiver')
subject = data.get('subject')
content = data.get('content')
if not all([receiver, subject, content]):
return jsonify({'error': '参数不完整'}), 400
# 提交异步任务,返回任务对象
task = send_email_task.delay(receiver, subject, content)
return jsonify({
'task_id': task.id,
'message': '发信任务已提交到后台执行'
}), 202
查询任务执行结果与状态
通过任务ID可以查询任务的执行状态、进度以及最终结果,Celery提供了对应的方法获取这些信息。
@app.route('/task_result/<task_id>', methods=['GET'])
def get_task_result(task_id):
"""查询任务结果的接口"""
task = send_email_task.AsyncResult(task_id)
if task.state == 'PENDING':
# 任务还未开始执行
response = {
'state': task.state,
'status': '任务等待执行中'
}
elif task.state == 'PROGRESS':
# 任务执行中
response = {
'state': task.state,
'progress': task.info.get('progress', 0),
'status': '任务执行中'
}
elif task.state == 'SUCCESS':
# 任务执行成功
response = {
'state': task.state,
'result': task.result,
'status': '任务执行成功'
}
else:
# 任务执行失败
response = {
'state': task.state,
'error': str(task.info),
'status': '任务执行失败'
}
return jsonify(response)
启动服务与测试验证
需要分别启动Flask应用和Celery worker进程,Celery worker会监听消息队列,自动获取并执行异步任务。
启动Flask应用
python app.py
启动Celery worker
在项目目录下执行以下命令,启动Celery worker进程,注意替换app为你的文件名。
celery -A app.celery worker --loglevel=info
测试流程
首先调用发信接口提交任务,拿到返回的task_id,再通过任务结果接口查询该task_id对应的任务状态,验证异步执行和结果查询功能是否正常。
注意事项
- 生产环境中建议给Celery配置独立的进程管理,比如使用supervisor管理Celery worker进程,避免进程意外退出。
- 消息队列和结果存储的后端地址需要根据实际部署环境调整,不要直接使用示例中的本地地址。
- 如果任务需要访问Flask应用的配置或者数据库,需要在任务内部手动推送Flask应用上下文,避免上下文丢失导致的问题。