生产环境偶现SQLAlchemy InterfaceError:无法在手动启动的事务中调用Connection.transaction()问题排查求助
生产环境偶现SQLAlchemy InterfaceError:无法在手动启动的事务中调用Connection.transaction()问题排查求助
问题背景
最近在生产环境遇到了偶发的SQLAlchemy报错,但开发环境完全正常,错误信息如下:
InterfaceError: (sqlalchemy.dialects.postgresql.asyncpg.InterfaceError) <class 'asyncpg.exceptions._base.InterfaceError'>: cannot use Connection.transaction() in a manually started transaction
[SQL: SELECT tags.id, tags.name, tags.color, tags.channel_id, tags.creator_id, tags.tag_type_id, tags.created_at, tags.updated_at
FROM tags
WHERE tags.channel_id = $1::BIGINT ORDER BY tags.created_at DESC]
[parameters: (-1001943261480,)]
这个错误明确提示:代码尝试在一个已经手动启动的事务中,再次触发事务相关操作导致冲突。结合我项目的代码片段,想请教大家生产环境中可能的触发原因是什么?
项目相关代码片段
IoC容器(Dishka Provider)
from typing import AsyncIterable from aiokafka import AIOKafkaConsumer 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.clients.api.invite_links import HelperBotCli from infrastructure.db.sqlalchemy.setup import create_engine, create_session_pool from infrastructure.message_brokers.kafka import KafkaMessageBroker 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()) @provide(scope=Scope.REQUEST) def helper_bot_client(self, config: Config) -> HelperBotCli: return HelperBotCli( base_url=config.helper_bot.HELPER_BOT_BASE_URL, create_url=config.helper_bot.ADD_LINK_ENDPOINT, delete_url=config.helper_bot.DELETE_LINK_ENDPOINT, ) @provide(scope=Scope.APP) async def kafka_broker(self, config: Config) -> AsyncIterable[KafkaMessageBroker]: consumer = AIOKafkaConsumer( config.rp_config.CONSUME_TOPIC, config.rp_config.CONSUME_MERGE_ACCOUNTS_TOPIC, bootstrap_servers=config.rp_config.RP_SERVER, group_id=config.rp_config.CONSUME_GROUP, enable_auto_commit=False, auto_offset_reset="earliest", sasl_mechanism=config.rp_config.SASL_MECHANISM, security_protocol=config.rp_config.SECURITY_PROTOCOL, sasl_plain_username=config.rp_config.RP_USER, sasl_plain_password=config.rp_config.RP_PASS, max_poll_records=config.rp_config.BATCH_SIZE, ) broker = KafkaMessageBroker(consumer) await broker.start() yield broker await broker.close()
SQLAlchemy初始化代码(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
工作单元(UoW)代码
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")
我的排查思路与疑问
- 生产环境流量远大于开发环境,会不会是并发/批量场景下的Session复用或事务嵌套触发了这个问题?
- 注意到Kafka消费者是
Scope.APP级别,但UoW是Scope.REQUEST级别——在Kafka批量消费消息时,Dishka的Request Scope会不会出现Session共享、事务状态混乱的情况? - 虽然注释了UoW中的
await self._session.begin(),但SQLAlchemy的AsyncSession在执行查询时,会自动开启隐式事务。会不会是代码中某个地方手动启动了事务,后续查询又触发了隐式事务导致冲突?
想请教各位大佬,生产环境中可能的根因是什么,以及如何验证和修复?
内容来源于stack exchange
相关产品推荐
相关产品推荐

