mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 01:10:23 +08:00
perf(stability): batch committer 异步化及降级冷却参数调优
- 将 batch_committer 的 DB commit 操作通过 asyncio.to_thread 移至线程池,避免阻塞事件循环 - 将 f-string 日志替换为 loguru 惰性格式化 - 降低事件循环延迟降级的冷却时间和乘数,加快降级恢复
This commit is contained in:
@@ -32,7 +32,7 @@ class BatchCommitter:
|
|||||||
"""启动后台批量提交任务"""
|
"""启动后台批量提交任务"""
|
||||||
if self._task is None:
|
if self._task is None:
|
||||||
self._task = asyncio.create_task(self._batch_commit_loop())
|
self._task = asyncio.create_task(self._batch_commit_loop())
|
||||||
logger.info(f"批量提交器已启动,间隔: {self.interval_seconds}s")
|
logger.info("批量提交器已启动,间隔: {}s", self.interval_seconds)
|
||||||
|
|
||||||
async def stop(self) -> Any:
|
async def stop(self) -> Any:
|
||||||
"""停止后台任务"""
|
"""停止后台任务"""
|
||||||
@@ -65,7 +65,7 @@ class BatchCommitter:
|
|||||||
await self._commit_all()
|
await self._commit_all()
|
||||||
raise
|
raise
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"批量提交出错: {e}")
|
logger.error("批量提交出错: {}", e)
|
||||||
|
|
||||||
async def _commit_all(self) -> None:
|
async def _commit_all(self) -> None:
|
||||||
"""提交所有待处理的 Session"""
|
"""提交所有待处理的 Session"""
|
||||||
@@ -79,22 +79,31 @@ class BatchCommitter:
|
|||||||
committed = 0
|
committed = 0
|
||||||
failed = 0
|
failed = 0
|
||||||
|
|
||||||
|
def _sync_commit_all() -> tuple[int, int]:
|
||||||
|
ok = 0
|
||||||
|
err = 0
|
||||||
for session in sessions_to_commit:
|
for session in sessions_to_commit:
|
||||||
try:
|
try:
|
||||||
session.commit()
|
session.commit()
|
||||||
committed += 1
|
ok += 1
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"提交 Session 失败: {e}")
|
logger.error("提交 Session 失败: {}", e)
|
||||||
try:
|
try:
|
||||||
session.rollback()
|
session.rollback()
|
||||||
except:
|
except:
|
||||||
pass
|
pass
|
||||||
failed += 1
|
err += 1
|
||||||
|
return ok, err
|
||||||
|
|
||||||
|
try:
|
||||||
|
committed, failed = await asyncio.to_thread(_sync_commit_all)
|
||||||
|
except Exception as e:
|
||||||
|
logger.error("批量提交线程异常: {}", e)
|
||||||
|
|
||||||
if committed > 0:
|
if committed > 0:
|
||||||
logger.debug(f"批量提交完成: {committed} 个 Session")
|
logger.debug("批量提交完成: {} 个 Session", committed)
|
||||||
if failed > 0:
|
if failed > 0:
|
||||||
logger.warning(f"批量提交失败: {failed} 个 Session")
|
logger.warning("批量提交失败: {} 个 Session", failed)
|
||||||
|
|
||||||
|
|
||||||
# 全局单例
|
# 全局单例
|
||||||
|
|||||||
@@ -34,8 +34,8 @@ _HEARTBEAT_DEDUP_TTL_SECONDS = 600
|
|||||||
_LOOP_WATCHDOG_INTERVAL_SECONDS = 1.0
|
_LOOP_WATCHDOG_INTERVAL_SECONDS = 1.0
|
||||||
_LOOP_LAG_WARNING_SECONDS = 1.0
|
_LOOP_LAG_WARNING_SECONDS = 1.0
|
||||||
_LOOP_LAG_DEGRADE_SECONDS = 3.0
|
_LOOP_LAG_DEGRADE_SECONDS = 3.0
|
||||||
_LOOP_LAG_DEGRADE_MIN_COOLDOWN_SECONDS = 10.0
|
_LOOP_LAG_DEGRADE_MIN_COOLDOWN_SECONDS = 2.0
|
||||||
_LOOP_LAG_DEGRADE_MAX_COOLDOWN_SECONDS = 30.0
|
_LOOP_LAG_DEGRADE_MAX_COOLDOWN_SECONDS = 5.0
|
||||||
_LOOP_LAG_WARNING_LOG_INTERVAL_SECONDS = 10.0
|
_LOOP_LAG_WARNING_LOG_INTERVAL_SECONDS = 10.0
|
||||||
|
|
||||||
_HOP_BY_HOP_HEADERS = frozenset(
|
_HOP_BY_HOP_HEADERS = frozenset(
|
||||||
@@ -168,7 +168,7 @@ class HubConnectionManager:
|
|||||||
if lag_seconds >= _LOOP_LAG_DEGRADE_SECONDS:
|
if lag_seconds >= _LOOP_LAG_DEGRADE_SECONDS:
|
||||||
cooldown = min(
|
cooldown = min(
|
||||||
_LOOP_LAG_DEGRADE_MAX_COOLDOWN_SECONDS,
|
_LOOP_LAG_DEGRADE_MAX_COOLDOWN_SECONDS,
|
||||||
max(_LOOP_LAG_DEGRADE_MIN_COOLDOWN_SECONDS, lag_seconds * 3.0),
|
max(_LOOP_LAG_DEGRADE_MIN_COOLDOWN_SECONDS, lag_seconds * 1.5),
|
||||||
)
|
)
|
||||||
degraded_until = now + cooldown
|
degraded_until = now + cooldown
|
||||||
self._degraded_until = max(self._degraded_until, degraded_until)
|
self._degraded_until = max(self._degraded_until, degraded_until)
|
||||||
|
|||||||
Reference in New Issue
Block a user