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")
其他可选方案:乐观锁
如果你的场景是高并发读写,行级锁可能会导致性能瓶颈,推荐使用乐观锁方案:
- 在
Task模型中添加一个版本号字段,比如version: int = Column(Integer, default=1) - 更新时通过版本号判断记录是否被修改:
@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
相关产品推荐
相关产品推荐

