From ce5c705ab40664dc8629b478e8933a6bb80a5366 Mon Sep 17 00:00:00 2001 From: shy3130 <415333856@qq.com> Date: Wed, 1 Jul 2026 11:19:32 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E8=B7=A8=E6=97=A5=E5=90=8E=E8=BF=9E?= =?UTF-8?q?=E6=9D=BF=E6=A2=AF=E9=98=9F=E6=95=B4=E4=BD=93=E5=B0=91=E7=AE=97?= =?UTF-8?q?=E4=B8=80=E6=A1=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 根因: 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 不反映到内存缓存) --- backend/app/jobs/daily_pipeline.py | 13 +++++++++---- backend/app/tickflow/repository.py | 31 +++++++++++++++++++++++++++++- 2 files changed, 39 insertions(+), 5 deletions(-) diff --git a/backend/app/jobs/daily_pipeline.py b/backend/app/jobs/daily_pipeline.py index 6bb7099..a8ad488 100644 --- a/backend/app/jobs/daily_pipeline.py +++ b/backend/app/jobs/daily_pipeline.py @@ -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"), diff --git a/backend/app/tickflow/repository.py b/backend/app/tickflow/repository.py index f916c40..7ddd521 100644 --- a/backend/app/tickflow/repository.py +++ b/backend/app/tickflow/repository.py @@ -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