73 lines
2.4 KiB
Python
73 lines
2.4 KiB
Python
"""四库统一注册表:异步生命周期编排。
|
|||
|
|
|
||
|
|
启动:逐库 init,失败按指数退避重试;DB_STRICT_STARTUP=true 时重试耗尽直接抛错阻止启动。
|
||
|
|
关闭:逆序 dispose。健康:check_ready_detail 返回带耗时明细,供 /health 与监控面板使用。
|
||
|
|
"""
|
||
|
|
import asyncio
|
||
|
|
import logging
|
||
|
|
import time
|
||
|
|
|
||
|
|
from . import milvus, mysql, neo4j, redis
|
||
|
|
from config.settings import settings
|
||
|
|
|
||
|
|
logger = logging.getLogger("config.database")
|
||
|
|
|
||
|
|
# 关闭顺序与设计一致:mysql → redis → neo4j → milvus
|
||
|
|
_DB_MODULES = (mysql, redis, neo4j, milvus)
|
||
|
|
|
||
|
|
|
||
|
|
def _name(m) -> str:
|
||
|
|
return m.__name__.rsplit(".", 1)[-1]
|
||
|
|
|
||
|
|
|
||
|
|
async def _init_with_retry(m, strict: bool) -> bool:
|
||
|
|
name = _name(m)
|
||
|
|
retries = settings.db.conn_retries
|
||
|
|
for attempt in range(retries):
|
||
|
|
try:
|
||
|
|
await m.init_db()
|
||
|
|
return True
|
||
|
|
except Exception as e:
|
||
|
|
logger.warning(
|
||
|
|
"[%s] init failed (attempt %d/%d): %s: %s",
|
||
|
|
name, attempt + 1, retries, type(e).__name__, e,
|
||
|
|
)
|
||
|
|
if attempt < retries - 1:
|
||
|
|
await asyncio.sleep(settings.db.retry_backoff_sec * (2**attempt))
|
||
|
|
if strict:
|
||
|
|
raise RuntimeError(f"[{name}] 初始化失败({retries} 次尝试后)")
|
||
|
|
return False
|
||
|
|
|
||
|
|
|
||
|
|
async def init_db() -> None:
|
||
|
|
"""手动预热(应用启动不自动调用,会话懒创建):asyncio.run(database.init_db())"""
|
||
|
|
for m in _DB_MODULES:
|
||
|
|
await _init_with_retry(m, settings.db.strict_startup)
|
||
|
|
|
||
|
|
|
||
|
|
async def dispose() -> None:
|
||
|
|
for m in _DB_MODULES:
|
||
|
|
try:
|
||
|
|
await m.dispose()
|
||
|
|
except Exception as e:
|
||
|
|
logger.warning("[%s] dispose failed: %s: %s", _name(m), type(e).__name__, e)
|
||
|
|
|
||
|
|
|
||
|
|
async def check_ready_detail() -> dict[str, dict]:
|
||
|
|
"""含耗时明细:{db: {status, ms}}。任一失败不影响其它库探测。"""
|
||
|
|
out: dict[str, dict] = {}
|
||
|
|
for m in _DB_MODULES:
|
||
|
|
name = _name(m)
|
||
|
|
t0 = time.perf_counter()
|
||
|
|
try:
|
||
|
|
await m.check_health()
|
||
|
|
status = "ok"
|
||
|
|
except Exception as e:
|
||
|
|
status = f"down: {type(e).__name__}"
|
||
|
|
out[name] = {"status": status, "ms": round((time.perf_counter() - t0) * 1000, 1)}
|
||
|
|
return out
|
||
|
|
|
||
|
|
|
||
|
|
async def check_ready() -> dict[str, str]:
|
||
|
|
"""兼容旧契约:{db: status_str}。"""
|
||
|
|
return {k: v["status"] for k, v in (await check_ready_detail()).items()}
|