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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 11:40:11