生产环境偶发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
相关产品推荐
相关产品推荐

