fix(db): 每请求独立 AsyncSession(根除并发 500)
- 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 <noreply@anthropic.com>
This commit is contained in:
+12
-2
@@ -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:
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user