Files
dpb/backend/app/core/database.py
T
34047007@qq.com 0b50e0adfa checkpoint: SQLite 清理——业务/测试层统一 PostgreSQL(删 aiosqlite/fakeredis)
- setting: DATABASE_TYPE 收窄为 postgres,DB_URI/ASYNC_DB_URI 去 sqlite 分支
- database: create_async_engine_and_session 去 sqlite 分支(同步 psycopg 引擎保留给 jobstore)
- number_gen: 去 DATABASE_TYPE 非 postgres 早退(advisory lock 恒定走 PG)
- chat/crud: 删 SqliteDb 分支(agno 无硬依赖)
- pyproject/requirements/uv.lock: 删 aiosqlite、fakeredis;conftest 口径注释同步
- test_stat_analysis_tc: docstring 口径 SQLite→PG
2026-08-07 21:47:42 +08:00

163 lines
5.5 KiB
Python

from fastapi import FastAPI
from redis import exceptions
from redis.asyncio import Redis
from sqlalchemy import Engine, create_engine
from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession, async_sessionmaker, create_async_engine
from sqlalchemy.orm import sessionmaker
from app.config.setting import settings
from app.core.base_model import MappedBase
from app.core.logger import logger
def create_engine_and_session(db_url: str = settings.DB_URI) -> tuple[Engine, sessionmaker]:
"""创建同步数据库引擎和会话工厂。
仅 APScheduler SQLAlchemyJobStore(同步库接口,apscheduler==3.11.0)消费;
业务层一律使用异步 create_async_engine_and_session,不得再引入同步引擎。
参数:
- db_url (str): 数据库连接URL,默认从配置中获取。
返回:
- tuple[Engine, sessionmaker]: 同步数据库引擎和会话工厂。
"""
try:
# 同步数据库引擎
engine: Engine = create_engine(
url=db_url,
echo=settings.DATABASE_ECHO,
pool_pre_ping=settings.POOL_PRE_PING,
pool_recycle=settings.POOL_RECYCLE,
)
except Exception as e:
logger.error(f"❌ 数据库连接失败 {e}")
raise
else:
# 同步数据库会话工厂
SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine)
return engine, SessionLocal
def create_async_engine_and_session(db_url: str = settings.ASYNC_DB_URI) -> tuple[AsyncEngine, async_sessionmaker[AsyncSession]]:
"""获取异步数据库会话连接。
参数:
- db_url (str): 异步数据库 URL,默认取配置项 ASYNC_DB_URI。
返回:
- tuple[AsyncEngine, async_sessionmaker[AsyncSession]]: 异步数据库引擎和会话工厂。
"""
try:
# 异步数据库引擎(统一 PostgreSQL 连接池参数)
async_engine = create_async_engine(
url=db_url,
echo=settings.DATABASE_ECHO,
echo_pool=settings.ECHO_POOL,
pool_pre_ping=settings.POOL_PRE_PING,
future=settings.FUTURE,
pool_recycle=settings.POOL_RECYCLE,
pool_size=settings.POOL_SIZE,
max_overflow=settings.MAX_OVERFLOW,
pool_timeout=settings.POOL_TIMEOUT,
pool_use_lifo=settings.POOL_USE_LIFO,
)
except Exception as e:
logger.error(f"❌ 数据库连接失败 {e}")
raise
else:
# 异步数据库会话工厂
AsyncSessionLocal = async_sessionmaker[AsyncSession](
bind=async_engine,
autocommit=settings.AUTOCOMMIT,
autoflush=settings.AUTOFLUSH if settings.AUTOFETCH is None else settings.AUTOFETCH,
expire_on_commit=settings.EXPIRE_ON_COMMIT,
class_=AsyncSession,
)
return async_engine, AsyncSessionLocal
# 同步引擎/会话仅 APScheduler SQLAlchemyJobStore 消费(同步库接口);业务层统一用 async_engine
engine, db_session = create_engine_and_session()
async_engine, async_db_session = create_async_engine_and_session()
async def check_db() -> None:
"""检查数据库连接是否正常。"""
try:
async with async_engine.connect():
pass
logger.info("✅ 数据库连接正常")
except Exception as e:
logger.error(f"❌ 数据库连接失败: {e}")
raise e
async def create_tables() -> None:
"""创建数据库表(根据 ORM metadata)。
返回:
- None
"""
try:
async with async_engine.begin() as coon:
await coon.run_sync(MappedBase.metadata.create_all)
except Exception as e:
logger.error(f"❌ 数据库表结构初始化失败: {e}")
raise e
async def drop_tables() -> None:
"""删除数据库表(根据 ORM metadata)。
返回:
- None
"""
try:
async with async_engine.begin() as conn:
await conn.run_sync(MappedBase.metadata.drop_all)
except Exception as e:
logger.error(f"❌ 数据库表结构删除失败: {e}")
raise e
async def redis_connect(app: FastAPI, status: bool) -> Redis | None:
"""创建或关闭Redis连接。
参数:
- app (FastAPI): FastAPI应用实例。
- status (bool): 连接状态,True为创建连接,False为关闭连接。
返回:
- Redis | None: Redis连接实例,如果连接失败则返回None。
"""
if status:
app.state.redis = None # 先置不可用;连接失败时应用以降级模式启动,不再中断
try:
rd = await Redis.from_url(
url=settings.REDIS_URI,
encoding="utf-8",
decode_responses=True,
protocol=2,
health_check_interval=settings.REDIS_HEALTH_CHECK_INTERVAL,
max_connections=settings.POOL_SIZE,
socket_timeout=settings.POOL_TIMEOUT,
)
if await rd.ping(): # pyright: ignore[reportGeneralTypeIssues]
app.state.redis = rd
return rd
logger.error("❌ 数据库 Redis ping 失败")
await rd.aclose()
except (exceptions.RedisError, OSError) as e:
logger.error(f"❌ 数据库 Redis 连接失败: {e}")
return None
else:
rd = app.state.redis
if rd is not None:
try:
await rd.close()
except Exception: # noqa: BLE001
pass
app.state.redis = None
logger.info("✅️ Redis连接已关闭")