FastAPI+SQLAlchemy+Asyncpg下PostgreSQL竞态条件问题求助
解决FastAPI+SQLAlchemy Async+PostgreSQL中的竞态插入问题
核心问题是先查询后插入的非原子操作导致并发场景下重复插入:多请求同时通过“记录不存在”的检查,最终触发重复数据报错。原代码中的with_for_update仅对已存在的行生效,无法锁定“不存在的行”,因此无法解决竞态。
以下是三层递进的有效解决方案:
1. 数据库层面添加联合唯一约束(必须)
首先在MedicoModel中定义联合唯一约束,让数据库从根源兜底禁止重复数据:
from sqlalchemy import Column, Integer, String, UniqueConstraint from sqlalchemy.ext.declarative import declarative_base Base = declarative_base() class MedicoModel(Base): __tablename__ = 'medicos' id = Column(Integer, primary_key=True, index=True) nome = Column(String, nullable=False) numero_registro = Column(String, nullable=False) sigla_conselho = Column(String, nullable=False) uf = Column(String, nullable=False) # 定义联合唯一约束,三个字段组合必须唯一 __table_args__ = ( UniqueConstraint( 'numero_registro', 'sigla_conselho', 'uf', name='uq_medico_registro_conselho_uf' ), )
2. 替换先查后插为原子性的INSERT ... ON CONFLICT操作
用PostgreSQL原生的INSERT ... ON CONFLICT DO NOTHING语法,把“检查存在+插入”合并为一个原子操作,彻底消除竞态。修改busca_cria_medico函数:
from sqlalchemy import insert, select async def busca_cria_medico(self, venda_medico: dict) -> MedicoModel: """原子化检查并创建Medico,避免竞态""" # 构造插入语句:冲突时不执行任何操作,返回插入的记录 insert_stmt = insert(MedicoModel).values( nome=venda_medico.nome, numero_registro=venda_medico.numero_registro, uf=venda_medico.uf, sigla_conselho=venda_medico.sigla_conselho ).on_conflict_do_nothing( index_elements=['numero_registro', 'sigla_conselho', 'uf'] ).returning(MedicoModel) result = await self.db.execute(insert_stmt) novo_medico = result.scalar_one_or_none() if novo_medico: await self.db.flush() return novo_medico # 未插入成功则直接查询已存在的记录返回 medico_existente = await self.db.execute( select(MedicoModel) .filter( MedicoModel.numero_registro == venda_medico.numero_registro, MedicoModel.sigla_conselho == venda_medico.sigla_conselho, MedicoModel.uf == venda_medico.uf ) ) return medico_existente.scalar_one()
为什么有效?
INSERT ... ON CONFLICT是数据库层面的原子操作,同一事务内完成“冲突检查+插入”,不会被其他并发请求打断。returning子句直接返回插入结果,避免额外查询;冲突时返回None,再查询已存在记录即可。
3. 异常兜底处理(可选但推荐)
极端情况下仍可能因事务隔离级别或其他原因触发IntegrityError,添加异常捕获确保逻辑健壮:
from sqlalchemy import insert, select from sqlalchemy.exc import IntegrityError async def busca_cria_medico(self, venda_medico: dict) -> MedicoModel: try: insert_stmt = insert(MedicoModel).values( nome=venda_medico.nome, numero_registro=venda_medico.numero_registro, uf=venda_medico.uf, sigla_conselho=venda_medico.sigla_conselho ).on_conflict_do_nothing( index_elements=['numero_registro', 'sigla_conselho', 'uf'] ).returning(MedicoModel) result = await self.db.execute(insert_stmt) novo_medico = result.scalar_one_or_none() if novo_medico: await self.db.flush() return novo_medico medico_existente = await self.db.execute( select(MedicoModel) .filter( MedicoModel.numero_registro == venda_medico.numero_registro, MedicoModel.sigla_conselho == venda_medico.sigla_conselho, MedicoModel.uf == venda_medico.uf ) ) return medico_existente.scalar_one() except IntegrityError: # 捕获唯一约束冲突,回滚子事务后查询已存在记录 await self.db.rollback() medico_existente = await self.db.execute( select(MedicoModel) .filter( MedicoModel.numero_registro == venda_medico.numero_registro, MedicoModel.sigla_conselho == venda_medico.sigla_conselho, MedicoModel.uf == venda_medico.uf ) ) return medico_existente.scalar_one()
额外优化建议
- 移除原代码中的
with_for_update:它仅对已存在行加锁,对“不存在的行”无效,反而增加不必要的锁开销。 - 保留当前的嵌套事务
async with self.db.begin_nested():子事务失败不会影响外层事务,适合此类独立操作场景。 - PostgreSQL默认的
READ COMMITTED隔离级别已足够处理该场景,无需调整。
内容的提问来源于stack exchange,提问作者João
相关产品推荐
相关产品推荐

