From 7ec9b8fd8a8c43fe643b3b74331a6b00ed18f748 Mon Sep 17 00:00:00 2001 From: Pine Date: Thu, 27 Aug 2026 14:17:55 +0800 Subject: [PATCH] =?UTF-8?q?fix(db):=20=E6=AF=8F=E8=AF=B7=E6=B1=82=E7=8B=AC?= =?UTF-8?q?=E7=AB=8B=20AsyncSession=EF=BC=88=E6=A0=B9=E9=99=A4=E5=B9=B6?= =?UTF-8?q?=E5=8F=91=20500=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - get_db 由「返回 app.state.db 单例」改为「每请求各建一个 Database 并在结束关闭」 - Database.__init__ 复用按 url 缓存的共享异步引擎(_shared_engine),避免每请求重建引擎 - 修复 GET /auth/me、/admin/* 并发时 sqlalchemy.exc.InvalidRequestError: concurrent operations are not permitted(同一 AsyncSession 被并发请求复用) - 对齐 park 子应用 _carrier_db 的每请求 Database 模式 Co-Authored-By: Claude --- app/api/dependencies.py | 14 ++++++++++++-- app/infrastructure/repositories.py | 24 ++++++++++++++++++++---- 2 files changed, 32 insertions(+), 6 deletions(-) diff --git a/app/api/dependencies.py b/app/api/dependencies.py index 772f6c7..bed83c2 100644 --- a/app/api/dependencies.py +++ b/app/api/dependencies.py @@ -23,8 +23,18 @@ def extract_bearer_token(request: Request) -> str: return "" -def get_db(request: Request) -> Database: - return request.app.state.db +async def get_db(): + """每请求一个独立 Database/AsyncSession(绑定共享引擎)。 + + 之前直接返回 app.state.db(单例 AsyncSession)被并发请求复用,会间歇报 + ``sqlalchemy.exc.InvalidRequestError: concurrent operations are not permitted``。 + 改为每请求各建一个 Database 并在请求结束关闭,根除该竞态。 + """ + db = Database() + try: + yield db + finally: + await db.close() async def _resolve_identity(user: dict, db: Database) -> dict: diff --git a/app/infrastructure/repositories.py b/app/infrastructure/repositories.py index f18bdc3..ab48849 100644 --- a/app/infrastructure/repositories.py +++ b/app/infrastructure/repositories.py @@ -2102,17 +2102,33 @@ class StatsRepository: # --------------------------------------------------------------------------- # 数据库门面 # --------------------------------------------------------------------------- +# 共享异步引擎(按 url 缓存,进程内复用)—— 每请求一个独立 Database/AsyncSession, +# 引擎跨请求共享,避免「同一 AsyncSession 被并发请求复用」的竞态。 +_SHARED_ENGINES: dict[str, "object"] = {} + + +def _shared_engine(url: str): + eng = _SHARED_ENGINES.get(url) + if eng is None: + from .db import make_async_engine + eng = make_async_engine(url) + _SHARED_ENGINES[url] = eng + return eng + + class Database: - """持有全部 Repository,统一访问入口(基础设施层)。""" + """持有全部 Repository,统一访问入口(基础设施层)。 + + ``session is None`` 时新建一个独立 AsyncSession(绑定共享引擎)—— 供每请求 + 各持一个 Database,避免并发复用同一 session 抛 ``concurrent operations are not permitted``。 + """ def __init__(self, db_url: str | None = None, session: AsyncSession | None = None): if session is None: from sqlalchemy.ext.asyncio import async_sessionmaker - from .db import make_async_engine - url = db_url or config.DATABASE_URL - self._engine = make_async_engine(url) + self._engine = _shared_engine(url) self._session_factory = async_sessionmaker( bind=self._engine, class_=AsyncSession, autoflush=False, expire_on_commit=False, )