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

FastAPI并发更新竞态条件错误:正确规避方法咨询

解决FastAPI接口并发更新数据库的竞态条件问题

原接口代码

你编写的FastAPI更新任务状态接口及重试函数如下:

@router.put("/{task_id}/status")
def update_status(task_id: str, task_data: UpdateStatus, db: Session = Depends(get_db)):
    task = retry_filter_by_id(db, model_name=Task, id=task_id)
    if task:
        task.external_status = task_data.status
        task.external_status_update_time = datetime.now()
        retry_commit(db)
        return JSONResponse(content={
            "status": "Updated",
            "task_id": task_id
        })
    raise HTTPException(status_code=400, detail=f"Task {task_id} not found")

# retry functions
@retry_db_operation(max_retries=3)
def retry_filter_by_id(db, model_name, id: str):
    return db.query(model_name).filter_by(id=id).one_or_none()

@retry_db_operation(max_retries=3)
def retry_commit(db):
    db.commit()

问题说明

当前接口在并发请求时出现Concurrent update error(并发更新错误),本质是竞态条件:多个请求同时读取同一任务记录,各自修改后提交,导致后提交的请求覆盖前一个的修改,或者触发数据库的并发冲突。

你的方案分析与修正

你提出的给查询添加行级锁的思路是正确的,with_for_update()会在查询时锁定目标记录,阻止其他事务修改该记录直到当前事务完成,从根源上避免竞态条件。不过你的代码存在一处小问题,同时需要在接口中正确调用带锁的查询:

修正后的重试查询函数

@retry_db_operation(max_retries=3)
def retry_filter_by_id(db, model_name, id: str, lock=False):
    query = db.query(model_name).filter_by(id=id)
    if lock:
        query = query.with_for_update()
    return query.one_or_none()

注意:原代码中最后一行one_or_none缺少括号,已修正为one_or_none()。

更新接口调用带锁查询

修改接口函数,调用retry_filter_by_id时传入lock=True,确保读取记录时就锁定:

@router.put("/{task_id}/status")
def update_status(task_id: str, task_data: UpdateStatus, db: Session = Depends(get_db)):
    # 读取时加行级锁
    task = retry_filter_by_id(db, model_name=Task, id=task_id, lock=True)
    if task:
        task.external_status = task_data.status
        task.external_status_update_time = datetime.now()
        retry_commit(db)
        return JSONResponse(content={
            "status": "Updated",
            "task_id": task_id
        })
    raise HTTPException(status_code=400, detail=f"Task {task_id} not found")

其他可选方案:乐观锁

如果你的场景是高并发读写,行级锁可能会导致性能瓶颈,推荐使用乐观锁方案:

  1. 在Task模型中添加一个版本号字段,比如version: int = Column(Integer, default=1)
  2. 更新时通过版本号判断记录是否被修改:
@retry_db_operation(max_retries=3)
def update_task_status(db, task_id: str, new_status: str):
    updated_rows = db.query(Task)\
        .filter(Task.id == task_id, Task.version == Task.version)\
        .update({
            "external_status": new_status,
            "external_status_update_time": datetime.now(),
            "version": Task.version + 1
        }, synchronize_session=False)
    db.commit()
    return updated_rows > 0

# 接口中调用
@router.put("/{task_id}/status")
def update_status(task_id: str, task_data: UpdateStatus, db: Session = Depends(get_db)):
    success = update_task_status(db, task_id, task_data.status)
    if success:
        return JSONResponse(content={
            "status": "Updated",
            "task_id": task_id
        })
    raise HTTPException(status_code=400, detail=f"Task {task_id} not found or concurrent update occurred")

乐观锁不需要加锁,通过版本号判断是否有并发修改,适合高并发场景,但需要处理更新失败的重试逻辑。

内容的提问来源于stack exchange,提问作者mascai

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.01 15:44:53