FastAPI+Celery环境下Update数据库语句执行无效问题
FastAPI + Celery + Repository模式下Update语句无生效问题
问题现象
在FastAPI项目中使用Repository模式结合Celery时,大部分数据库查询操作正常,但某条Update语句执行后,数据库对应字段未发生变化。日志显示已返回目标记录的ID,但数据库中reminder_sent_at字段仍为空。
相关代码与信息
Update查询语句
UPDATE_REMINDER_SENT_AT_QUERY = """ UPDATE public.trainings SET reminder_sent_at = :next_reminder_date AT TIME ZONE 'utc' WHERE id = :training_id RETURNING id; """
仓库层代码
async def update_reminder_sent_at( self, training_id: UUID, next_reminder_date ): async with self.db.transaction(): record: int = await self.db.fetch_val( query=UPDATE_REMINDER_SENT_AT_QUERY, values={ "training_id": training_id, "next_reminder_date": next_reminder_date, }, ) logger.info(f'record {record}')
服务层代码
async def update_reminder_sent_at( self, training_id: UUID, next_reminder_date ): await self.repository.update_reminder_sent_at(training_id, next_reminder_date)
Celery任务调用
@app_celery.task def task_update( training_id: UUID, user_id: UUID, tenant_id: UUID ): """Send email reminders for a specific training and user.""" async def send_reminder(training_id: UUID, user_id: UUID, tenant_id: UUID): next_reminder_date = datetime.utcnow() + timedelta( days=training.reminder.schedule_in_days ) await training_service.update_reminder_sent_at( training_id, next_reminder_date ) asyncio.run(send_reminder(training_id, user_id, tenant_id))
日志输出
celery_worker | 2024-06-15 12:00:00.388 |INFO | api.trainings.repository:update_reminder_sent_at:163 - record d1d33ada-1e98-40e7-982c-f553b09bcaa0
数据库查询结果
laas_api=# select id, reminder_sent_at from trainings laas_api-# ; id | reminder_sent_at --------------------------------------+------------------ b8b6b8d3-20ca-4ced-9720-4b6585f116d9 | 3f9e08af-dacf-4979-97be-bbe3f86c5979 | fba1a12d-a8d9-418d-851a-a55e29c02996 | 679596a0-264c-4dad-aca2-37c197534621 | 24d3469d-bfa7-43f0-95b9-996105af6df0 | 45213723-0711-4d86-97ef-d9c7528bcebd | d1d33ada-1e98-40e7-982c-f553b09bcaa0 | (7 rows)
数据库连接配置
@asynccontextmanager async def get_db(): database = Database( settings.database_url, force_rollback=True, min_size=3, max_size=20 ) await database.connect() try: yield database finally: await database.disconnect()
问题根源
数据库连接配置中的force_rollback=True是核心问题。这个参数会强制所有数据库事务在会话结束后自动回滚,哪怕代码里显式开启了事务并执行了Update操作,最终都不会提交到数据库。该参数一般仅用于测试环境,避免测试数据污染正式库,但如果在生产或正常运行环境启用,就会导致所有写操作失效。
修复方案
1. 调整数据库连接配置
移除force_rollback=True,或者根据运行环境动态设置(仅在测试环境启用):
@asynccontextmanager async def get_db(): database = Database( settings.database_url, # 仅在测试环境开启强制回滚 force_rollback=settings.ENVIRONMENT == "test", min_size=3, max_size=20 ) await database.connect() try: yield database finally: await database.disconnect()
2. 验证时间字段处理逻辑
检查next_reminder_date的格式与数据库reminder_sent_at字段类型是否匹配:
- 如果
reminder_sent_at是TIMESTAMP WITH TIME ZONE类型,且next_reminder_date已经是UTC时间,SQL语句中无需额外加AT TIME ZONE 'utc',可以简化为:
UPDATE_REMINDER_SENT_AT_QUERY = """ UPDATE public.trainings SET reminder_sent_at = :next_reminder_date WHERE id = :training_id RETURNING id; """
额外排查点
- 确认Celery任务中
training.reminder.schedule_in_days是有效数值,避免生成异常的next_reminder_date - 查看数据库系统日志,确认Update语句是否实际执行,是否存在隐式回滚的记录
- 验证仓库层事务逻辑,确保没有未捕获的异常导致事务自动回滚
内容的提问来源于stack exchange,提问作者Lutaaya Huzaifah Idris
相关产品推荐
相关产品推荐

