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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 05:55:00