"""一次性跑完所有 sync task(按依赖顺序)。 按 sort_order 触发:stock_basic → kline_daily → kline_index → kline_5min → industry_sector → sector_features → share_snapshot → market_regime 每个 task 跑完写日志 + 把 result 落到 logs/runall_summary.json。 注:以下 task **不**在 runall 链中 —— 各自有独立的 systemd timer: - tick_trade : mairui 21:00 发布 → market-sync-tick.timer (21:05) - moneyflow : mairui 21:30 发布 → market-sync-moneyflow.timer (21:35) """ from __future__ import annotations import json import sys import time from pathlib import Path _PROJECT_ROOT = Path(__file__).resolve().parent.parent if str(_PROJECT_ROOT) not in sys.path: sys.path.insert(0, str(_PROJECT_ROOT)) from app.core.utils.logging import setup_logging, get_logger setup_logging() logger = get_logger("runall") from app.core.datasource.registry import build_default_registry from app.core.sync.registry import seed_sync_registry build_default_registry() seed_sync_registry() from app.tasks import get_task # 顺序执行,按 sort_order 走(便于排查 + 避免外部 API 限流) # 注:tick_trade / moneyflow 由独立 timer 跑(21:05 / 21:35) TASKS = [ "stock_basic", "kline_index", "kline_daily", "kline_5min", "industry_sector", "sector_features", "share_snapshot", "market_regime", ] # 不跳任何 task — 全量跑 SKIP: set[str] = set() results = {} t_all = time.time() for tid in TASKS: if tid in SKIP: results[tid] = {"status": "skipped", "message": "已 ok"} logger.info(f"[runall] {tid} 跳过(已 ok)") continue logger.info(f"[runall] >>> 开始 {tid}") t0 = time.time() try: task = get_task(tid) # 大表 task 限小批量股票,避免跑爆;其他 task 跑全量 kwargs = {"max_workers": 5} if tid in {"kline_daily", "kline_5min", "moneyflow", "share_snapshot"}: # 这些 task 支持 MARKET_DATA_STOCK_LIMIT 环境变量,但通过 cli 跑可 --codes 限 # 测全量费时,先跑全量 pass r = task.run(trigger_source="runall", **kwargs) elapsed = round(time.time() - t0, 1) r["elapsed_sec"] = elapsed results[tid] = r logger.info(f"[runall] <<< {tid} 完成 status={r.get('status')} elapsed={elapsed}s msg={r.get('message')}") except Exception as e: elapsed = round(time.time() - t0, 1) logger.exception(f"[runall] !!! {tid} 异常: {e}") results[tid] = {"status": "error", "message": str(e), "elapsed_sec": elapsed} # 每个 task 之间 sleep 5s 让健康监控跑 time.sleep(5) elapsed_total = round(time.time() - t_all, 1) summary = {"total_elapsed_sec": elapsed_total, "tasks": results} out = _PROJECT_ROOT / "logs" / "runall_summary.json" out.write_text(json.dumps(summary, ensure_ascii=False, indent=2)) logger.info(f"[runall] 全部完成 total={elapsed_total}s summary={out}") print(json.dumps(summary, ensure_ascii=False, indent=2))