【问题标题】:How do I manually commit a sqlalchemy database transaction inside a pyramid web app?如何在金字塔 Web 应用程序中手动提交 sqlalchemy 数据库事务?
【发布时间】:2021-04-10 18:43:15
【问题描述】:

我有一个 Pyramid Web 应用程序,在将更改提交到 sqlalchemy 数据库后需要运行 Celery 任务。我知道我可以使用 request.tm.get().addAfterCommitHook() 来做到这一点。但是,这对我不起作用,因为我还需要在视图中使用 celery 任务的 task_id。因此,我需要在对 Celery 任务调用 task.delay() 之前提交对数据库的更改。

zope.sqlalchemy 文档说我可以使用 transaction.commit() 手动提交。但是,这对我不起作用; celery 任务在更改提交到数据库之前运行,即使我在调用 task.delay() 之前调用了 transaction.commit()

我的 Pyramid 视图代码如下所示:

ride=appstruct_to_ride(dbsession,appstruct)
dbsession.add(ride)

# Flush dbsession so ride gets an id assignment
dbsession.flush()

# Store ride id
ride_id=ride.id
log.info('Created ride {}'.format(ride_id))

# Commit ride to database
import transaction
transaction.commit()

# Queue a task to update ride's weather data
from ..processing.weather import update_ride_weather
update_weather_task=update_ride_weather.delay(ride_id)

url = self.request.route_url('rides')
return HTTPFound(
    url,
    content_type='application/json',
    charset='',
    text=json.dumps(
        {'ride_id':ride_id,
         'update_weather_task_id':update_weather_task.task_id}))

我的 celery 任务如下所示:

@celery.task(bind=True,ignore_result=False)
def update_ride_weather(self,ride_id, train_model=True):

    from ..celery import session_factory
    
    logger.debug('Received update weather task for ride {}'.format(ride_id))

    dbsession=session_factory()
    dbsession.expire_on_commit=False

    with transaction.manager:
        ride=dbsession.query(Ride).filter(Ride.id==ride_id).one()

celery 任务因 NoResultFound 失败:

  File "/app/cycling_data/processing/weather.py", line 478, in update_ride_weather
    ride=dbsession.query(Ride).filter(Ride.id==ride_id).one()
  File "/usr/local/lib/python3.8/site-packages/sqlalchemy/orm/query.py", line 3282, in one
    raise orm_exc.NoResultFound("No row was found for one()")

当我事后检查数据库时,我看到记录实际上是在 celery 任务运行并失败之后创建的。所以这意味着 transaction.commit() 没有按预期提交事务,而是在视图返回后由 zope.sqlalchemy 机器自动提交更改。如何在我的视图代码中手动提交事务?

【问题讨论】:

    标签: python sqlalchemy celery pyramid


    【解决方案1】:

    request.tmpyramid_tm 定义,可以是线程本地的transaction.manager 对象或按请求对象,具体取决于您如何配置pyramid_tm(查找pyramid_tm.manager_hook 在某处定义以确定哪个一个正在使用中。

    您的问题很棘手,因为您所做的任何事情都应该适合pyramid_tm 以及它期望的操作方式。具体来说,它计划在请求的生命周期内控制事务——提早提交对于该事务不是一个好主意。 pyramid_tm 试图帮助提供故障保护功能,以在请求生命周期中的任何位置发生任何故障时回滚整个请求 - 而不仅仅是在您的视图中可调用。

    选项 1:

    无论如何,请尽早提交。如果您要执行此操作,则提交后的失败无法回滚已提交的数据,因此您可能会部分提交请求。好的,很好,这是你的问题,所以答案是使用request.tm.commit() 可能后跟request.tm.begin() 为任何后续更改开始一个新的。您还需要小心不要跨该边界共享 sqlalchemy 托管对象,例如 request.user 等,因为它们需要刷新/合并到新事务中(SQLAlchemy 的身份缓存默认情况下不能信任从不同事务加载的数据因为这就是隔离级别的工作方式)。

    选项 2:

    为您想提早提交的数据启动一个单独的事务。好的,所以假设你没有使用像transaction.managerscoped_session 这样的任何threadlocals,那么你可能可以开始自己的事务并提交它,而无需触及由pyramid_tm 控制的dbsession。一些适用于 pyramid-cookiecutter-starter 项目结构的通用代码可能是:

    from myapp.models import get_tm_session
    
    tmp_tm = transaction.TransactionManager(explicit=True)
    with tmp_tm:
        dbsession_factory = request.registry['dbsession_factory']
        tmp_dbsession = get_tm_session(dbsession_factory, tmp_tm)
        # ... do stuff with tmp_dbsession that is committed in this with-statement
        ride = appstruct_to_ride(tmp_dbsession, appstruct)
        # do not use this ride object outside of the with-statement
        tmp_dbsession.add(ride)
        tmp_dbsession.flush()
        ride_id = ride.id
    
    # we are now committed so go ahead and start your background worker
    update_weather_task = update_ride_weather.delay(ride_id)
    
    # maybe you want the ride object outside of the tmp_dbsession
    ride = dbsession.query(Ride).filter(Ride.id==ride_id).one()
    
    return {...}
    

    这还不错 - 可能是在不将 celery 挂接到 pyramid_tm 控制的 dbsession 的情况下,就故障模式而言,你能做的最好的事情。

    【讨论】:

    • 谢谢@michael。我发现您的回答有助于我更好地理解交易机制。但是,我的 celery 任务仍然使用这两种方法引发了 NoResultFound。如果我在提交事务后将dbsession.query(Ride).filter(Ride.id==ride_id).one() 插入到视图中,我也会在那里得到一个 NoResultFound。
    • 上述评论的例外是,在我的单元测试(使用 sqlite 数据库)中运行查询 dbsession.query(Ride).filter(Ride.id==ride_id).one() 工作正常,但在模拟生产环境(使用 mariadb 数据库)中相同查询失败。和以前一样,事务在视图返回后被提交,但手动提交似乎对 mariadb 没有任何作用。
    • 顺便说一句,我查看了我的 app/models/__init__.py,它有 settings['tm.manager_hook'] = 'pyramid_tm.explicit_manager',根据 pyramid documentation 是一个非线程本地管理器。
    • 我粘贴的代码 sn-p,如果放入可调用的视图中,应该可以工作。如果在提交之后调用 celery 任务,如在粘贴中,并且它返回 NoResultFound,我怀疑在我粘贴的 sn-p 中有一些东西不能正常工作,你使用 dbsession 而不是 tmp_dbsession也许吧。
    • 看来我的问题可能是我在粘贴代码之前使用 request.dbsession 运行查询。似乎如果我在创建 tmp_dbsession 之前对 request.dbsession 运行任何查询(即使它只是一个 SELECT),那么我会收到 NoResultFound 错误。
    猜你喜欢
    • 1970-01-01
    • 2020-07-06
    • 2013-04-26
    • 2014-10-05
    • 2023-04-02
    • 1970-01-01
    • 2013-11-23
    • 2013-01-27
    • 1970-01-01
    相关资源
    最近更新 更多