【发布时间】: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