diff --git a/scripts/data_platform/backfill_15min_baostock.py b/scripts/data_platform/backfill_15min_baostock.py index 891d667..d3e026a 100644 --- a/scripts/data_platform/backfill_15min_baostock.py +++ b/scripts/data_platform/backfill_15min_baostock.py @@ -22,6 +22,7 @@ BaoStock特点: 变更记录: v1.0 (2026-05-05) 赵云创建 v1.1 (2026-05-05) 修复:旧新浪数据含后复权值,改为全量重建模式 + v1.2 (2026-07-13) 硬化:异常强制重登 / login SIGALRM 超时 / 断路器 / marker-based resume """ import argparse @@ -29,6 +30,7 @@ import json import logging import os import shutil +import signal import sys import time from datetime import datetime @@ -50,8 +52,10 @@ ALL_STOCKS_FILE = NAS_ROOT / "A股数据" / "stock_info" / "stock_basic_info_raw # BaoStock配置 BS_START_DATE = os.environ.get("BS_START_DATE", "2024-01-01") BS_INTERVAL = 0.4 # 请求间隔秒 -BS_MAX_RETRIES = 2 +BS_MAX_RETRIES = 3 PROGRESS_SAVE_EVERY = 500 +BS_LOGIN_TIMEOUT = 30 # bs.login() 超时秒数(防永久挂起) +CIRCUIT_BREAKER_THRESHOLD = 30 # 连续失败N只 → 断路退出 # ======================== 日志 ======================== @@ -100,16 +104,29 @@ def is_backfilled(parquet_path: Path) -> bool: def load_progress() -> set: - progress_file = PROGRESS_DIR / "backfill_15min_progress.json" - if progress_file.exists(): - try: - return set(json.loads(progress_file.read_text()).get("done", [])) - except Exception: - pass - return set() + """ + 从 marker 文件构造已完成集合(真相源)。 + 只有 .{prefix}{code}_15min.baostock marker 存在 = 该票已成功回填。 + failed 票无 marker,下次复跑会被重试,避免缺口永不重试。 + JSON 进度文件只用于观察,不再决定 todo。 + """ + done = set() + if not MINUTE_15_DIR.exists(): + return done + suffix = "_15min.baostock" + for marker in MINUTE_15_DIR.glob(f".*{suffix}"): + name = marker.name # .sh600000_15min.baostock + stripped = name[1:-len(suffix)] # sh600000 + if len(stripped) < 8: # sh/sz(2) + 6位代码 + continue + code = stripped[2:] # 600000 + if len(code) == 6 and code.isdigit(): + done.add(code) + return done def save_progress(done_set: set): + """写 JSON 进度文件(仅观察用,load_progress 不再读取)""" PROGRESS_DIR.mkdir(parents=True, exist_ok=True) progress_file = PROGRESS_DIR / "backfill_15min_progress.json" progress_file.write_text(json.dumps({ @@ -118,6 +135,61 @@ def save_progress(done_set: set): }, ensure_ascii=False)) +# ======================== BaoStock 连接管理 ======================== + +class _LoginTimeout(Exception): + """bs.login() SIGALRM 超时""" + + +def _login_timeout_handler(signum, frame): + raise _LoginTimeout("bs.login() 超时") + + +def _login_with_timeout(timeout: int = BS_LOGIN_TIMEOUT) -> bool: + """ + 给 bs.login() 套 SIGALRM 超时,防止卡死时整个回填永久挂起。 + signal.alarm 只能在主线程用(本脚本主线程,OK)。 + 返回 True=登录成功,False=失败/超时。 + """ + old_handler = signal.signal(signal.SIGALRM, _login_timeout_handler) + signal.alarm(timeout) + try: + lg = bs.login() + ok = lg.error_code == "0" + if not ok: + logger.error("bs.login() 失败: %s", lg.error_msg) + return ok + except _LoginTimeout: + logger.error("bs.login() 超时 %d 秒(baostock 疑似限流冷却中)", timeout) + return False + except Exception as e: + logger.error("bs.login() 异常: %s", e) + return False + finally: + signal.alarm(0) + signal.signal(signal.SIGALRM, old_handler) + + +def _relogin() -> bool: + """ + 强制重登:logout + login(强制建新 socket)。 + 登录失败再试 1 次。baostock 是模块级单连接,本脚本单进程串行,重登安全。 + """ + try: + bs.logout() + except Exception: + pass + if _login_with_timeout(): + return True + # 第一次失败,等2秒再试一次 + time.sleep(2) + try: + bs.logout() + except Exception: + pass + return _login_with_timeout() + + # ======================== 核心逻辑 ======================== def fetch_bs_15min(bs_code: str, start_date: str, end_date: str) -> Optional[pd.DataFrame]: @@ -186,9 +258,17 @@ def backfill_one(code: str, start_date: str, end_date: str, force: bool = False) df_new = fetch_bs_15min(bs_code, start_date, end_date) if df_new is not None and len(df_new) > 0: break + # 空数据:可能是真的无数据(新股/停牌),也可能是查询异常 + logger.debug("backfill %s 空数据重试 %d/%d", code, attempt + 1, BS_MAX_RETRIES) except Exception as e: - logger.debug("backfill %s 重试%d: %s", code, attempt + 1, e) - time.sleep(1) + # socket 死了(如 [Errno 57] Socket is not connected),同 socket 重试必然再败 + # 强制重登建新 socket 再 retry + logger.warning("backfill %s 异常重试 %d/%d: %s — 强制重登", + code, attempt + 1, BS_MAX_RETRIES, e) + if not _relogin(): + logger.error("重登失败,放弃 %s", code) + df_new = None + break if df_new is None or df_new.empty: return "failed", 0 @@ -233,10 +313,9 @@ def main(): logger.error("❌ NAS未挂载") sys.exit(1) - # 登录BaoStock - lg = bs.login() - if lg.error_code != "0": - logger.error("❌ BaoStock登录失败: %s", lg.error_msg) + # 登录BaoStock(带超时,防永久挂起) + if not _login_with_timeout(): + logger.error("❌ BaoStock登录失败(baostock 疑似冷却中),请稍后重试") sys.exit(1) logger.info("✅ BaoStock登录成功") @@ -274,17 +353,16 @@ def main(): stats = {"ok": 0, "skipped": 0, "failed": 0, "rows": 0} t_start = time.time() RELOGIN_EVERY = 400 # 每400只重新登录BaoStock,防止连接断开 + consecutive_failures = 0 # 连续失败计数(断路器):ok 重置,failed +1,skipped 中性 + circuit_triggered = False for i, code in enumerate(todo): - # 定期重新登录保持连接 + # 定期重新登录保持连接(用带超时的 _relogin,防永久挂起) if i > 0 and i % RELOGIN_EVERY == 0: - bs.logout() - time.sleep(2) - lg = bs.login() - if lg.error_code != "0": - logger.error("BaoStock重连失败: %s,等待30秒", lg.error_msg) + if not _relogin(): + logger.error("BaoStock 定期重连失败,等待30秒后再试") time.sleep(30) - lg = bs.login() + _relogin() logger.info("BaoStock重连 @ %d/%d", i, len(todo)) try: @@ -296,9 +374,27 @@ def main(): stats[status] = stats.get(status, 0) + 1 if status == "ok": stats["rows"] += total_rows + consecutive_failures = 0 + elif status == "failed": + consecutive_failures += 1 + # skipped 中性:不重置(不证明 baostock 可用),不递增(不算失败) done_set.add(code) + # 断路器:连续N只全 failed → baostock 疑似不可达,保存进度主动退出 + if consecutive_failures >= CIRCUIT_BREAKER_THRESHOLD: + logger.error( + "❌ 断路器触发:连续 %d 只失败,baostock 疑似不可达," + "保存进度后退出(等冷却后复跑,done_set 不含 failed 票会被重试)", + consecutive_failures) + save_progress(done_set) + try: + bs.logout() + except Exception: + pass + circuit_triggered = True + break + if (i + 1) % PROGRESS_SAVE_EVERY == 0: save_progress(done_set) elapsed = time.time() - t_start @@ -313,13 +409,23 @@ def main(): # 保存最终进度 save_progress(done_set) - bs.logout() + try: + bs.logout() + except Exception: + pass elapsed = time.time() - t_start logger.info("=" * 60) - logger.info("✅ 回补完成,耗时 %.1f 秒", elapsed) + if circuit_triggered: + logger.info("❌ 回补因断路器中止,耗时 %.1f 秒", elapsed) + else: + logger.info("✅ 回补完成,耗时 %.1f 秒", elapsed) logger.info("统计: %s", json.dumps(stats, ensure_ascii=False)) + # 断路器触发 → 非零退出码(cron/调度器可识别) + if circuit_triggered: + sys.exit(2) + if __name__ == "__main__": main() diff --git a/scripts/data_platform/resume_5yr_watcher.py b/scripts/data_platform/resume_5yr_watcher.py new file mode 100644 index 0000000..97c975a --- /dev/null +++ b/scripts/data_platform/resume_5yr_watcher.py @@ -0,0 +1,68 @@ +"""5yr 15min 回填冷却后自动续跑 watcher。 + +baostock 被限流/冷却时登录会卡死或失败。本 watcher 每 30min 探测一次 baostock, +可登就启动 backfill_15min_baostock(marker-based resume 自动续跑缺口),断路器触发 +或登录失败就继续等。backfill 成功(exit 0,全部缺口填完)则退出。 + + detached 启动: + nohup ./venv311/bin/python3 scripts/data_platform/resume_5yr_watcher.py \ + > /Users/chufeng/data_cache/stock/logs/daily_update/resume_watcher.log 2>&1 &! +""" +import os +import subprocess +import time + +ROOT = "/Users/chufeng/.openclaw/sanguo_projects/sanguo_vnpy_v2" +PY = os.path.join(ROOT, "venv311/bin/python3") +SCRIPT = os.path.join(ROOT, "scripts/data_platform/backfill_15min_baostock.py") +PROBE_INTERVAL = 1800 # 30min + +# 探测脚本:SIGALRM 给 bs.login() 套 15s 闹钟,避免卡死 +_PROBE = r''' +import signal +def _h(*a): raise TimeoutError() +signal.signal(signal.SIGALRM, _h) +signal.alarm(15) +try: + import baostock as bs + lg = bs.login() + ok = lg.error_code == "0" + try: bs.logout() + except Exception: pass + print("LOGIN_OK" if ok else "LOGIN_FAIL") +except Exception as e: + print("LOGIN_ERR", type(e).__name__, e) +finally: + signal.alarm(0) +''' + + +def baostock_up() -> bool: + try: + r = subprocess.run([PY, "-c", _PROBE], capture_output=True, text=True, timeout=30) + return "LOGIN_OK" in r.stdout + except Exception: + return False + + +def main() -> None: + round_ = 0 + while True: + round_ += 1 + if baostock_up(): + print(f"[round {round_}] baostock 可登,启动回填(marker-resume 续跑缺口)", flush=True) + env = dict(os.environ) + env["STOCK_ROOT"] = "/Users/chufeng/data_cache/stock" + env["BS_START_DATE"] = "20210101" + rc = subprocess.call([PY, SCRIPT], env=env, cwd=ROOT) + if rc == 0: + print(f"[round {round_}] 回填完成(exit 0),退出 watcher", flush=True) + return + print(f"[round {round_}] 回填 exit={rc}(断路器/失败),等 {PROBE_INTERVAL}s 再试", flush=True) + else: + print(f"[round {round_}] baostock 不可登(冷却中),等 {PROBE_INTERVAL}s", flush=True) + time.sleep(PROBE_INTERVAL) + + +if __name__ == "__main__": + main()