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

生产环境偶发sqlalchemy.InterfaceError:手动启动事务后无法调用Connection.transaction()排查求助

生产环境偶发sqlalchemy.InterfaceError:手动启动事务后无法调用Connection.transaction()排查求助

大家好,我遇到了一个非常棘手的偶发问题,想请各位帮忙排查一下:

我们的代码在多个项目中都稳定运行,但其中一个服务在生产环境偶尔会抛出这个错误:

sqlalchemy.exc.InterfaceError: (sqlalchemy.dialects.postgresql.asyncpg.InterfaceError) <class 'asyncpg.exceptions._base.InterfaceError'>: cannot use Connection.transaction() in a manually started transaction

奇怪的是,这个问题在Dev环境和本地完全复现不出来。我已经排查过代码,确认除了IoC容器和消费者里的事务管理外,没有其他手动操作事务的代码。

下面是相关的代码片段,麻烦各位帮忙看看可能哪里出了问题:

1. IoC容器 Provider 代码

from dishka import Provider, Scope, provide
from sqlalchemy.ext.asyncio import AsyncEngine
from sqlalchemy.orm import sessionmaker
from config import Config, get_config
from infrastructure.db.sqlalchemy.setup import create_engine, create_session_pool
from infrastructure.uow.sqlalchemy import SQLAlchemyUnitOfWork

class BaseProvider(Provider):
    @provide(scope=Scope.APP)
    def config(self) -> Config:
        return get_config()

    @provide(scope=Scope.APP)
    def sqlalchemy_engine(self, config: Config) -> AsyncEngine:
        return create_engine(config.pg_config)

    @provide(scope=Scope.APP)
    def session_pool(self, sqlalchemy_engine: AsyncEngine) -> sessionmaker:
        return create_session_pool(sqlalchemy_engine)

    @provide(scope=Scope.REQUEST)
    def sqlalchemy_uow(self, session_factory: sessionmaker) -> SQLAlchemyUnitOfWork:
        return SQLAlchemyUnitOfWork(_session=session_factory())

2. 数据库初始化 setup.py

from typing import Callable, AsyncContextManager
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine, AsyncEngine
from sqlalchemy.orm import sessionmaker
from config import PGConfig
from infrastructure.db.sqlalchemy.models import BaseModel

def create_engine(db: PGConfig, echo: bool = False) -> AsyncEngine:
    engine = create_async_engine(
        db.pg_database_url,
        echo=echo
    )
    return engine

async def create_tables(engine: AsyncEngine) -> None:
    async with engine.begin() as conn:
        await conn.run_sync(BaseModel.metadata.create_all)

def create_session_pool(engine: AsyncEngine) -> Callable[[], AsyncContextManager[AsyncSession]]:
    session_pool = sessionmaker(bind=engine, expire_on_commit=False, class_=AsyncSession)
    return session_pool

3. 工作单元(UoW)代码

为了尝试解决问题,我已经注释掉了await self._session.begin(),但还是偶发错误:

from dataclasses import dataclass
from types import TracebackType
from typing import Self
from sqlalchemy.ext.asyncio import AsyncSession
from infrastructure.repositories.invite_links.sqlalchemy import SQLAlchemyInviteLinksRepository
from infrastructure.repositories.tags.sqlalchemy import SQLAlchemyTagsRepository, SQLAlchemyTagsTypesRepository
from infrastructure.uow.base import BaseUnitOfWork

@dataclass
class SQLAlchemyUnitOfWork(BaseUnitOfWork):
    _session: AsyncSession = None

    async def __aenter__(self) -> Self:
        # await self._session.begin()
        return self

    async def __aexit__(
        self, exc_type: type[BaseException] | None, exc_val: BaseException | None, exc_tb: TracebackType | None
    ) -> None:
        try:
            if exc_type is None:
                await self._session.commit()
            else:
                await self._session.rollback()
        finally:
            await self._session.close()

    @property
    def invite_links(self) -> SQLAlchemyInviteLinksRepository:
        self._session_check()
        return SQLAlchemyInviteLinksRepository(self._session)

    @property
    def tags(self) -> SQLAlchemyTagsRepository:
        self._session_check()
        return SQLAlchemyTagsRepository(self._session)

    @property
    def tags_types(self) -> SQLAlchemyTagsTypesRepository:
        self._session_check()
        return SQLAlchemyTagsTypesRepository(self._session)

    def _session_check(self) -> None:
        if not self._session:
            raise ValueError("Session is not initialized")

4. 消费者相关的 lifespan.py

我把session factory传入了kafka消费者:

from contextlib import asynccontextmanager
import aiojobs
from fastapi import FastAPI
from sqlalchemy.ext.asyncio import AsyncEngine
from sqlalchemy.orm import sessionmaker
from app.api.background import consume_in_background
from config import Config
from infrastructure.db.sqlalchemy.setup import create_tables
from infrastructure.message_brokers.kafka import KafkaMessageBroker

@asynccontextmanager
async def lifespan(app: FastAPI) -> FastAPI:
    container = app.state.dishka_container
    engine = await container.get(AsyncEngine)
    scheduler = aiojobs.Scheduler()
    broker = await container.get(KafkaMessageBroker)
    config = await container.get(Config)

    if config.app.is_dev or config.app.is_production:
        session_factory = await container.get(sessionmaker)
        job = await scheduler.spawn(consume_in_background(broker, session_factory, config))

    await create_tables(engine)
    yield

    if config.app.is_dev or config.app.is_production:
        await job.close()

我的排查思路和疑问:

  • 已经确认代码中没有手动调用begin()开启额外事务的地方,UoW里也注释掉了session的begin
  • 问题只在生产环境偶发,Dev和本地完全正常,会不会是生产环境的并发量、数据库连接池配置的问题?
  • 消费者中使用session factory的方式有没有问题?会不会在消费消息时,session的事务状态被意外改变了?

有没有遇到过类似问题的朋友,或者能从代码中发现潜在问题的?麻烦指点一下,非常感谢!

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 10:28:07