You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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官方装饰器(推荐)

  1. 修复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
  1. 修改/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": "用户创建成功"}

额外注意事项

  1. 确保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,
    # 其他环境配置...
}
  1. 启动Celery Worker时指定应用实例:
celery -A app.celery worker --loglevel=info
  1. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 23:05:24