提交
This commit is contained in:
@@ -339,6 +339,13 @@ async def _run_sync(days: int, force: bool) -> None:
|
||||
|
||||
# candles/复权因子已更新:作废旧 K 线预览缓存(键含版本号,自增即全体失效)
|
||||
await cache.bump_version("candles")
|
||||
# 预热同步状态缓存:同步任务自己付一次重聚合(>10s),轮询方毫秒级拿到新数字
|
||||
_sync_state["step"] = "正在更新统计缓存"
|
||||
try:
|
||||
async with async_session() as s2:
|
||||
await _db_stats(s2)
|
||||
except Exception: # noqa: BLE001 —— 预热失败只影响下一次轮询的时延
|
||||
pass
|
||||
_sync_state["step"] = "同步完成"
|
||||
except Exception as e: # noqa: BLE001
|
||||
_sync_state["error"] = f"同步失败:{str(e)[:300]}"
|
||||
@@ -362,31 +369,56 @@ async def start_sync(session: AsyncSession, days: int, force: bool) -> dict:
|
||||
return dict(_sync_state)
|
||||
|
||||
|
||||
# candles 是千万行表,count 较重;前端每 2s 轮询状态,需 TTL 缓存降载
|
||||
_status_stats_cache: dict = {"at": 0.0, "data": None}
|
||||
_STATS_TTL = 30.0
|
||||
# candles 是千万行表,重聚合(全表 count / distinct 日期)在远程库实测 >11s;
|
||||
# 结果按 ver:candles 版本号缓存(同步完成即 bump 失效),前端 2s 轮询只付毫秒级。
|
||||
_status_stats_cache: dict = {"at": 0.0, "ver": -1, "data": None}
|
||||
_STATS_TTL = 120.0 # 进程内兜底 TTL(Redis 不可用时重聚合的最小间隔)
|
||||
_stats_bg_tasks: set[asyncio.Task] = set() # 后台写缓存的引用,防 GC
|
||||
_stats_lock = asyncio.Lock() # 单飞锁:同步尾部的预热与轮询并发时,重聚合只跑一次
|
||||
|
||||
|
||||
async def _db_stats(session: AsyncSession) -> dict:
|
||||
"""candles/快照/股票列表实况(30s TTL 缓存)。"""
|
||||
now = time.time()
|
||||
if _status_stats_cache["data"] is not None and now - _status_stats_cache["at"] < _STATS_TTL:
|
||||
return _status_stats_cache["data"]
|
||||
async def _heavy_stats(session: AsyncSession) -> dict:
|
||||
"""重聚合:行数/日期数。4 条查询走千万行表(>11s),绝不能落在轮询热路径上。"""
|
||||
stocks = int(await session.scalar(select(func.count()).select_from(StockBasic)) or 0)
|
||||
daily_rows = int(await session.scalar(select(func.count()).select_from(Candle)) or 0)
|
||||
snap_rows = int(await session.scalar(select(func.count()).select_from(DailySnapshot)) or 0)
|
||||
last_daily = await session.scalar(
|
||||
select(func.max(Candle.ts)).where(Candle.timeframe == "1d")
|
||||
)
|
||||
n_dates = int(await session.scalar(
|
||||
select(func.count(func.distinct(func.date(Candle.ts)))).where(Candle.timeframe == "1d")
|
||||
) or 0)
|
||||
data = {
|
||||
"stocks": stocks, "daily_rows": daily_rows, "snapshot_rows": snap_rows,
|
||||
"last_daily": last_daily, "dates": n_dates,
|
||||
return {
|
||||
"stocks": stocks, "daily_rows": daily_rows,
|
||||
"snapshot_rows": snap_rows, "dates": n_dates,
|
||||
}
|
||||
_status_stats_cache.update(at=now, data=data)
|
||||
return data
|
||||
|
||||
|
||||
async def _db_stats(session: AsyncSession) -> dict:
|
||||
"""candles/快照/股票列表实况:重聚合走「版本化 Redis + 进程内」双层缓存,
|
||||
last_daily(max(ts),走索引很快)保持每次实时——它是 UI 主展示字段。"""
|
||||
|
||||
def _fresh_local(ver: int) -> dict | None:
|
||||
d = _status_stats_cache["data"]
|
||||
if d is not None and _status_stats_cache["ver"] == ver \
|
||||
and time.time() - _status_stats_cache["at"] < _STATS_TTL:
|
||||
return d
|
||||
return None
|
||||
|
||||
ver = await cache.get_version("candles")
|
||||
heavy = _fresh_local(ver)
|
||||
if heavy is None:
|
||||
async with _stats_lock: # 双检:等锁期间可能已被并发请求/同步预热填充
|
||||
heavy = _fresh_local(ver) or await cache.cache_get(f"syncstats:v{ver}")
|
||||
if heavy is None:
|
||||
heavy = await _heavy_stats(session)
|
||||
# 写 Redis 后台执行,失败由 cache 层静默降级,不拖慢本次返回
|
||||
task = asyncio.create_task(
|
||||
cache.cache_set(f"syncstats:v{ver}", heavy, ttl=settings.sync_stats_redis_ttl))
|
||||
_stats_bg_tasks.add(task)
|
||||
task.add_done_callback(_stats_bg_tasks.discard)
|
||||
_status_stats_cache.update(at=time.time(), ver=ver, data=heavy)
|
||||
last_daily = await session.scalar(
|
||||
select(func.max(Candle.ts)).where(Candle.timeframe == "1d")
|
||||
)
|
||||
return {**heavy, "last_daily": last_daily}
|
||||
|
||||
|
||||
async def get_sync_status(session: AsyncSession) -> dict:
|
||||
|
||||
Reference in New Issue
Block a user