如何在Celery与SQLAlchemy之间保留Flask应用上下文?
Flask+Celery任务中SQLAlchemy无应用上下文报错
应用结构
|-app.py |-project |-__init__.py |-celery_utils.py |-config.py |-users |-__init__.py |-models.py |-tasks.py
相关代码
app.py
from project import create_app, ext_celery app = create_app() celery = ext_celery.celery @app.route("/") def alive(): return "alive"
/project/init.py
import os from flask import Flask from flask_celeryext import FlaskCeleryExt from flask_migrate import Migrate from flask_sqlalchemy import SQLAlchemy from project.celery_utils import make_celery from project.config import config # 实例化扩展 db = SQLAlchemy() migrate = Migrate() ext_celery = FlaskCeleryExt(create_celery_app=make_celery) def create_app(config_name=None): if config_name is None: config_name = os.environ.get("FLASK_CONFIG", "development") # 实例化Flask应用 app = Flask(__name__) # 加载配置 app.config.from_object(config[config_name]) # 初始化扩展 db.init_app(app) migrate.init_app(app, db) ext_celery.init_app(app) # 注册蓝图 from project.users import users_blueprint app.register_blueprint(users_blueprint) # Flask CLI shell上下文 @app.shell_context_processor def ctx(): return {"app": app, "db": db} return app
/project/celery_utils.py
from celery import current_app as current_celery_app def make_celery(app): celery = current_celery_app celery.config_from_object(app.config, namespace="CELERY") return celery
/project/users/init.py
from flask import Blueprint, request, jsonify from celery.result import AsyncResult from .tasks import post_to_db users_blueprint = Blueprint("users", __name__, url_prefix="/users", template_folder="templates") from . import models, tasks @users_blueprint.route('/users', methods=['POST']) def users(): request_data = request.get_json() task = post_to_db.delay(request_data) response = { "id": task.task_id, "status": task.status } return jsonify(response) @users_blueprint.route('/responses', methods=['GET']) def responses(): request_data = request.get_json() result = AsyncResult(id=request_data['id']) response = result.get() return jsonify(response)
/project/users/models.py
from project import db class User(db.Model): """用户模型""" __tablename__ = "users" id = db.Column(db.Integer, primary_key=True, autoincrement=True) username = db.Column(db.String(128), unique=True, nullable=False) email = db.Column(db.String(128), unique=True, nullable=False) def __init__(self, username, email, *args, **kwargs): self.username = username self.email = email
/project/users/tasks.py
from celery import shared_task from .models import User from project import db @shared_task() def post_to_db(payload): print("任务执行中") user = User(**payload) db.session.add(user) db.session.commit() db.session.close() return True
报错信息
RuntimeError: No application found. Either work inside a view function or push an application context. ...
已尝试的解决方法
- 使用
ext_celery装饰任务,触发报错:TypeError: exceptions must derive from BaseException - 在任务中手动推送
app.app_context(),触发相同TypeError - 从app中导入Celery实例装饰任务,仍触发相同TypeError
解决方案
核心问题
Celery Worker运行在独立进程中,未自动加载Flask应用上下文,导致SQLAlchemy无法关联应用配置。以下是两种可行修复方案:
方案1:任务内手动创建并推送应用上下文
修改/project/users/tasks.py,在任务内部初始化Flask应用并推送上下文:
from celery import shared_task from project import create_app, db from .models import User @shared_task() def post_to_db(payload): # 创建应用实例并推送上下文 app = create_app() with app.app_context(): user = User(**payload) db.session.add(user) db.session.commit() return {"success": True, "message": "用户创建成功"}
方案2:使用Flask-CeleryExt官方装饰器(推荐)
- 修复
celery_utils.py,确保正确创建Celery实例并自动发现任务:
from celery import Celery def make_celery(app): celery = Celery(app.import_name) celery.config_from_object(app.config, namespace="CELERY") # 自动发现users模块下的任务 celery.autodiscover_tasks(['project.users']) return celery
- 修改
/project/users/tasks.py,用ext_celery.task替代shared_task:
from project import ext_celery, db from .models import User @ext_celery.task() def post_to_db(payload): user = User(**payload) db.session.add(user) db.session.commit() return {"success": True, "message": "用户创建成功"}
额外注意事项
- 确保
config.py配置Celery必需参数:
class DevelopmentConfig: # Flask配置 DEBUG = True SQLALCHEMY_DATABASE_URI = 'sqlite:///dev.db' SQLALCHEMY_TRACK_MODIFICATIONS = False # Celery配置 CELERY_BROKER_URL = 'redis://localhost:6379/0' CELERY_RESULT_BACKEND = 'redis://localhost:6379/0' config = { 'development': DevelopmentConfig, # 其他环境配置... }
- 启动Celery Worker时指定应用实例:
celery -A app.celery worker --loglevel=info
responses路由中需从已初始化的Celery实例获取AsyncResult,修改如下:
@users_blueprint.route('/responses', methods=['GET']) def responses(): request_data = request.get_json() from project import ext_celery result = ext_celery.celery.AsyncResult(id=request_data['id']) response = {"status": result.status, "result": result.get()} return jsonify(response)
内容的提问来源于stack exchange,提问作者741852963
相关产品推荐
相关产品推荐

