fix: 跨日后连板梯队整体少算一档

根因: live_agg 缓存的昨日连板数 (_prev_consec_up/down) 跨日不失效,
盘中增量计算永远基于前天基准, 导致次日开盘连板数整体 -1。

- get_live_agg() 加跨日失效守卫: date.today() 变化时按需重建
  (fast-path 仅比较 today ~0.24ms/次, 跨日当天一次性磁盘校验)
- cron 盘后管道补 refresh_cache(), 与手动触发 /api/pipeline/run 对齐
  (此前 cron 路径漏了这步, 盘后落盘的新 enriched 不反映到内存缓存)
This commit is contained in:
shy3130
2026-07-01 11:19:32 +08:00
parent 4a7d405d9d
commit ce5c705ab4
2 changed files with 39 additions and 5 deletions
+9 -4
View File
@@ -772,11 +772,16 @@ def start_scheduler(repo: KlineRepository, capset: CapabilitySet) -> AsyncIOSche
)
# 盘后: 日 K + enriched(时间由偏好决定)
def _pipeline_then_refresh(on_progress=None):
# 与手动触发 (/api/pipeline/run) 对齐: 管道落盘后重建 Polars 内存缓存,
# 否则 live_agg 的昨日连板数等基准列会停留在旧交易日, 次日开盘连板梯队
# 整体少算一档 (仅手动触发或重启才会刷缓存, cron 调度路径此前漏了这步)。
result = run_now(repo, capset, on_progress=on_progress)
repo.refresh_cache()
return result
scheduler.add_job(
lambda: _run_tracked(
lambda on_progress=None: run_now(repo, capset, on_progress=on_progress),
"daily_pipeline",
),
lambda: _run_tracked(_pipeline_then_refresh, "daily_pipeline"),
trigger=CronTrigger(day_of_week="mon-fri",
hour=sched["hour"], minute=sched["minute"],
timezone="Asia/Shanghai"),
+30 -1
View File
@@ -282,6 +282,7 @@ class KlineRepository:
self._enriched_cache_date: date | None = None
self._live_agg_cache: pl.DataFrame | None = None # 预计算聚合表 (~5500行)
self._live_agg_cache_date: date | None = None
self._live_agg_check_date: date | None = None # 上次跨日校验时的 today (快路径节流)
self._instruments_cache: pl.DataFrame | None = None
# 完整 enriched 历史 (含所有指标, 供 filter_history 策略使用)
self._enriched_history_cache: pl.DataFrame | None = None # ~100万行
@@ -337,6 +338,7 @@ class KlineRepository:
self._enriched_history_start = None
self._live_agg_cache = None
self._live_agg_cache_date = None
self._live_agg_check_date = None
self._instruments_cache = None
self._index_instruments_cache = None
self._etf_enriched_cache = None
@@ -796,9 +798,36 @@ class KlineRepository:
return df.sort(["symbol", "date"])
def get_live_agg(self) -> pl.DataFrame:
"""返回盘中实时指标预计算聚合表。如无缓存则懒加载。"""
"""返回盘中实时指标预计算聚合表。如无缓存则懒加载。
live_agg 的核心列 _prev_consec_up/down (昨日连板数) 取自基准日 enriched。
基准日由 _live_agg_baseline_date 决定: 盘中(today 有实时分区) 取上一交易日,
非盘中(磁盘最新日 < today) 取该最新日本身。一旦跨日, 期望基准日会前移,
旧缓存会把连板数整体少算一档, 故这里除首次懒加载外还要校验基准日是否仍
符合当前预期, 不符则重建 (无需等盘后管道刷缓存)。
性能: get_live_agg 被每轮实时行情调用 (expert 档 1s 一次)。跨日只在
date.today() 翻天时发生, 故先用 today 做廉价的 fast-path (μs 级),
仅当 today 变化时才查磁盘确认 (DuckDB 扫 132 万行约 100ms+) 并按需重建。
"""
if self._live_agg_cache is None:
self._refresh_enriched()
self._live_agg_check_date = date.today() # 刚建过, 当天不必再查磁盘
else:
today = date.today()
if self._live_agg_check_date != today:
# today 翻天了 (次日开盘首次轮询): 校验基准日是否需要前移重建。
# 同一天内多次调用直接跳过, 避免每轮都扫 parquet。
self._live_agg_check_date = today
disk_latest = self._latest_enriched_date_duckdb()
if disk_latest is not None:
expected = self._live_agg_baseline_date(disk_latest)
if self._live_agg_cache_date != expected:
logger.info(
"live_agg 跨日失效, 重建: 缓存基准=%s, 期望基准=%s",
self._live_agg_cache_date, expected,
)
self._refresh_enriched()
if self._live_agg_cache is None:
return pl.DataFrame()
return self._live_agg_cache