【问题标题】:"ResourceClosedError: The transaction is closed" error with celery beat and sqlalchemy + pyramid app“ResourceClosedError:事务已关闭”错误与 celery beat 和 sqlalchemy + 金字塔应用程序
【发布时间】:2013-04-26 15:53:19
【问题描述】:

我有一个名为 mainsite 的金字塔应用程序。

网站以相当异步的方式工作,主要是通过从视图启动线程来执行后端操作。

它使用 sqlalchemy 连接到 mysql,并使用 ZopeTransactionExtension 进行会话管理。

到目前为止,该应用程序运行良好。

我需要在其上运行定期作业,并且它需要使用一些从视图启动的相同异步函数。

我使用了 apscheduler,但遇到了问题。于是想到了用celery beat作为一个单独的进程,把mainapp当作一个库,导入要使用的函数。

我的 celery 配置如下所示:

from datetime import timedelta
from api.apiconst import RERUN_CHECK_INTERVAL, AUTOMATION_CHECK_INTERVAL, \
    AUTH_DELETE_TIME

BROKER_URL = 'sqla+mysql://em:em@localhost/edgem'
CELERY_RESULT_BACKEND = "database"
CELERY_RESULT_DBURI = 'mysql://em:em@localhost/edgem'

CELERYBEAT_SCHEDULE = {
    'rerun': {
        'task': 'tasks.rerun_scheduler',
        'schedule': timedelta(seconds=RERUN_CHECK_INTERVAL)
    },
    'automate': {
        'task': 'tasks.automation_scheduler',
        'schedule': timedelta(seconds=20)
    },
    'remove-tokens': {
        'task': 'tasks.token_remover_scheduler',
        'schedule': timedelta(seconds=2 * 24 * 3600 )
    },
}

CELERY_TIMEZONE = 'UTC'

tasks.py 是

from celery import Celery
celery = Celery('tasks')
celery.config_from_object('celeryconfig')


@celery.task
def rerun_scheduler():
    from mainsite.task import check_update_rerun_tasks
    check_update_rerun_tasks()


@celery.task
def automation_scheduler():
    from mainsite.task import automate
    automate()


@celery.task
def token_remover_scheduler():
    from mainsite.auth_service import delete_old_tokens
    delete_old_tokens()

请记住,上述所有函数都会立即返回,但如果需要则启动线程

线程通过 transaction.commit() after session.add(object) 将对象保存到数据库中。

问题是整个事情像宝石一样只工作大约 30 分钟。之后ResourceClosedError: The transaction is closed 错误开始发生在有transaction.commit() 的地方。我不确定是什么问题,我需要帮助进行故障排除。

我在任务中导入的原因是为了摆脱这个错误。认为每次需要运行任务时都导入是个好主意,而且我可能每次都会得到一个新事务,但看起来情况并非如此。

【问题讨论】:

    标签: sqlalchemy celery pyramid


    【解决方案1】:

    根据我的经验,尝试将配置为与 Pyramid(使用 ZopeTransactionExtension 等)与 Celery 工作人员一起使用的会话重用会导致严重的难以调试的混乱。

    ZopeTransactionExtension 将 SQLAlchemy 会话绑定到 Pyramid 的请求 - 响应周期 - 事务自动启动并提交或回滚,您通常不应该在代码中使用 transaction.commit() - 如果一切正常,中兴通讯将提交所有内容,如果您的代码引发异常,您的事务将被回滚。

    使用 Celery,您需要手动管理 SQLAlchemy 会话,中兴通讯阻止您这样做,因此您需要以不同方式配置您的 DBSession

    像这样简单的东西会起作用:

    DBSession = None
    
    def set_dbsession(session):
        global DBSession
        if DBSession is not None:
            raise AttributeError("DBSession has been already set to %s!" % DBSession)
    
        DBSession = session
    

    然后从 Pyramid 启动代码你做

    def main(global_config, **settings):
        ...
        set_dbsession(scoped_session(sessionmaker(extension=ZopeTransactionExtension())))
    

    使用 Celery 有点棘手 - 我最终为 Celery 创建了一个自定义启动脚本,我在其中配置了会话。

    workersetup.py 蛋中:

      entry_points="""
      # -*- Entry points: -*-
      [console_scripts]
      custom_celery = worker.celeryd:start_celery
      custom_celerybeat = worker.celeryd:start_celerybeat
      """,
      )
    

    worker/celeryd.py:

    def initialize_async_session(db_string, db_echo):
    
        import sqlalchemy as sa
        from db import Base, set_dbsession
    
        session = sa.orm.scoped_session(sa.orm.sessionmaker(autoflush=True, autocommit=True))
        engine = sa.create_engine(db_string, echo=db_echo)
        session.configure(bind=engine)
    
        set_dbsession(session)
        Base.metadata.bind = engine
    
    
    def start_celery():
        initialize_async_session(DB_STRING, DB_ECHO)
        import celery.bin.celeryd
        celery.bin.celeryd.main()
    

    如果您打算将应用程序部署到生产服务器,那么您使用“从视图启动线程以执行后端操作”的一般方法对我来说有点危险 - Web 服务器经常回收,杀死或创建新的“工人”,因此通常不能保证每个特定进程都会在当前的请求-响应周期之后继续存在。不过我从来没有尝试过这样做,所以也许你会没事的:)

    【讨论】:

    • 嘿。我对这个男人感激不尽。我在上面坐了几天才想到在这里问。我会试试这个,解决方案似乎可行。我会试试这个,让你知道结果如何。
    猜你喜欢
    • 1970-01-01
    • 2021-04-10
    • 2013-11-23
    • 2014-01-15
    • 1970-01-01
    • 2015-09-17
    • 2017-06-11
    • 1970-01-01
    • 2014-04-20
    相关资源
    最近更新 更多