From 57f417e6eba11e05f7f6db16c36cb340373daecf Mon Sep 17 00:00:00 2001 From: wshy Date: Sat, 4 Jul 2026 15:53:50 +0800 Subject: [PATCH] =?UTF-8?q?fix(sync):=20=E4=BF=AE=E5=A4=8D=E5=90=8C?= =?UTF-8?q?=E6=AD=A5=E5=8D=A1=E6=AD=BB=E5=9C=A826%=E5=8F=8A=E5=90=AF?= =?UTF-8?q?=E5=8A=A8=E6=9C=9F=E8=A7=86=E5=9B=BE=E6=B3=A8=E5=86=8C=E5=B4=A9?= =?UTF-8?q?=E6=BA=83=20(#47)=20(#48)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 三处根因: 1. 限流逻辑失效 (kline_sync.py ×3, index_sync.py ×2) 条件 `len(chunks) > rpm` 在 free 档恒为 False (56 chunks < 60 rpm), 导致节流 sleep 永不执行, 请求密集打出触发服务端限流 → 卡死。 改为始终按 interval (60/rpm) 节流。 2. 卡死看门狗无法自愈 (pipeline_jobs.py, api/pipeline.py) 原超时检查只在再次点「同步」时触发, 等待中永不回收。 抽出 JobStore.reap_stale() (STALE_JOB_TIMEOUT_S=600), 在 /run 和 /jobs/{id} 轮询端点都调用 — 前端每秒轮询, 10分钟后 必定自动回收卡死 job, UI 不再永久停在 26%。 3. 启动期视图注册异常捕获不足 (repository.py) _register_views 只捕获 duckdb.IOException, 跨版本/平台空目录 可能抛 CatalogException 等炸掉 lifespan。放宽到 Exception。 Co-authored-by: shy3130 --- backend/app/api/pipeline.py | 22 +++++------------ backend/app/services/index_sync.py | 4 +-- backend/app/services/kline_sync.py | 6 ++--- backend/app/services/pipeline_jobs.py | 35 +++++++++++++++++++++++++++ backend/app/tickflow/repository.py | 5 +++- 5 files changed, 50 insertions(+), 22 deletions(-) diff --git a/backend/app/api/pipeline.py b/backend/app/api/pipeline.py index 11f522b..f39611d 100644 --- a/backend/app/api/pipeline.py +++ b/backend/app/api/pipeline.py @@ -29,22 +29,9 @@ async def run_now(request: Request) -> dict: repo = request.app.state.repo capset = request.app.state.capabilities - # 检测卡死的 running job (如 reload 后孤儿 task) - existing_id = job_store.active_id() - if existing_id: - existing = job_store.get(existing_id) - if existing and existing["status"] == "running": - from datetime import datetime, timezone - started = existing.get("started_at") - if started: - try: - start_dt = datetime.fromisoformat(started.replace("Z", "+00:00")) - elapsed = (datetime.now(timezone.utc) - start_dt).total_seconds() - if elapsed > 600: # 超过 10 分钟视为卡死 - logger.warning("强制取消卡死 job %s (已运行 %.0fs)", existing_id, elapsed) - job_store.fail(existing_id, "超时自动取消 (疑似 reload 后孤儿 task)") - except Exception: - pass + # 检测卡死的 running job (如 reload 后孤儿 task / 网络读无限阻塞)。 + # reap_stale 会在 /run 和 /jobs/{id} 轮询端点都调用,保证卡死后能自愈。 + job_store.reap_stale() job_id = job_store.create() @@ -81,6 +68,9 @@ async def run_now(request: Request) -> dict: @router.get("/jobs/{job_id}") def get_job(job_id: str) -> dict: + # 每次轮询都检查卡死 job — 前端每秒轮询,STALE_JOB_TIMEOUT_S(10min)后必定自愈, + # 无需用户再次手动点「同步」。 + job_store.reap_stale() j = job_store.get(job_id) if not j: raise HTTPException(status_code=404, detail="job not found") diff --git a/backend/app/services/index_sync.py b/backend/app/services/index_sync.py index 42f3b1b..1f474e3 100644 --- a/backend/app/services/index_sync.py +++ b/backend/app/services/index_sync.py @@ -233,7 +233,7 @@ def sync_and_persist_index_daily( interval = (60.0 / rpm) if rpm else 0 chunks = [symbols[i:i + batch_size] for i in range(0, len(symbols), batch_size)] for i, chunk in enumerate(chunks): - if i > 0 and interval > 0 and len(chunks) > rpm: + if i > 0 and interval > 0: import time time.sleep(interval) raw = kline_sync.sync_daily_batch( @@ -332,7 +332,7 @@ def sync_and_persist_etf_daily( chunks = [symbols[i:i + batch_size] for i in range(0, len(symbols), batch_size)] factors = _load_etf_factors(repo) for i, chunk in enumerate(chunks): - if i > 0 and interval > 0 and len(chunks) > rpm: + if i > 0 and interval > 0: import time time.sleep(interval) raw = kline_sync.sync_daily_batch( diff --git a/backend/app/services/kline_sync.py b/backend/app/services/kline_sync.py index 7a4d5f5..ba9a751 100644 --- a/backend/app/services/kline_sync.py +++ b/backend/app/services/kline_sync.py @@ -92,7 +92,7 @@ def sync_daily_batch(symbols: list[str], chunks = [symbols[i:i + batch_size] for i in range(0, len(symbols), batch_size)] for i, chunk in enumerate(chunks): - if i > 0 and interval > 0 and len(chunks) > rpm: + if i > 0 and interval > 0: time.sleep(interval) try: if start_time and end_time: @@ -295,7 +295,7 @@ def sync_adj_factor(symbols: list[str], repo: KlineRepository, all_dfs: list[pl.DataFrame] = [] for i, chunk in enumerate(chunks): - if i > 0 and interval > 0 and len(chunks) > rpm: + if i > 0 and interval > 0: time.sleep(interval) try: raw = tf.klines.ex_factors(chunk, **sdk_kwargs) @@ -427,7 +427,7 @@ def sync_minute_batch( chunks = [symbols[i:i + batch_size] for i in range(0, len(symbols), batch_size)] for i, chunk in enumerate(chunks): - if i > 0 and interval > 0 and len(chunks) > rpm: + if i > 0 and interval > 0: time.sleep(interval) try: if start_time and end_time: diff --git a/backend/app/services/pipeline_jobs.py b/backend/app/services/pipeline_jobs.py index 979eddb..81ec0c6 100644 --- a/backend/app/services/pipeline_jobs.py +++ b/backend/app/services/pipeline_jobs.py @@ -23,6 +23,11 @@ logger = logging.getLogger(__name__) JobStatus = Literal["pending", "running", "succeeded", "failed"] +# 运行超过此秒数视为卡死(reload 后孤儿 task / 网络读无限阻塞等)。 +# 由 reap_stale() 在 /run 和 /jobs/{id} 轮询端点检查 — 保证卡死后能自愈, +# 无需用户再次点击「同步」。 +STALE_JOB_TIMEOUT_S = 600 + def _default_store_dir() -> Path: from app.config import settings @@ -208,6 +213,36 @@ class JobStore: def active_id(self) -> str | None: return self._active_id + def reap_stale(self, timeout_s: int = STALE_JOB_TIMEOUT_S) -> None: + """回收运行超过 timeout_s 的卡死 running job(标记为 failed)。 + + 在 /run 和 /jobs/{id} 轮询端点都会调用 — 保证卡死后任意轮询都能自愈, + 无需用户再次手动触发同步。reload 后的孤儿 task(内存里已无 job 记录) + 不在此处理:它们没有 active_id,只能靠 executor 线程自然结束或进程重启。 + """ + with self._lock: + jid = self._active_id + if not jid: + return + j = self._active_jobs.get(jid) + if not j or j.get("status") != "running": + return + started = j.get("started_at") + if not started: + return + # 时间计算放到锁外(避免 datetime 解析持锁)。 + # started_at 形如 "2026-07-04T12:00:00Z"(start() 用 datetime.utcnow 存)。 + # 两端都用 timezone-aware UTC 比较,避免 naive/aware 混用导致 TypeError。 + try: + start_dt = datetime.fromisoformat(started.replace("Z", "+00:00")) + elapsed = (datetime.now(start_dt.tzinfo) - start_dt).total_seconds() + except Exception: # noqa: BLE001 + return + if elapsed > timeout_s: + logger.warning("reap_stale: 强制取消卡死 job %s (已运行 %.0fs)", + jid, elapsed) + self.fail(jid, f"超时自动取消 (运行 {int(elapsed)}s, 疑似卡死)") + def clear(self) -> None: """清空所有任务(内存 + 磁盘文件)。""" with self._lock: diff --git a/backend/app/tickflow/repository.py b/backend/app/tickflow/repository.py index beaaa7e..260acaf 100644 --- a/backend/app/tickflow/repository.py +++ b/backend/app/tickflow/repository.py @@ -180,7 +180,10 @@ class DataStore: for sql in statements: try: self.db.execute(sql) - except duckdb.IOException: + except Exception as e: # noqa: BLE001 + # 空数据目录(首次启动)或权限问题时 DuckDB 会抛 IOException; + # 跨版本/平台也可能抛 CatalogException 等。空目录缺视图不影响启动 + # (后续同步写入数据后会重新刷新视图),这里一律降级为 debug 日志。 logger.debug("view registration skipped (no parquet yet): %s", sql[:60]) self._register_unified_views()