fix(data): backfill硬化—重登on-error/登录超时/断路器/marker-resume+冷却续跑watcher

首跑5yr 673/5493后baostock连接死、定时重连没恢复、登录卡死=被限流。修四个硬伤:
- 重登on-error: except分支强制bs.logout+login建新socket再retry(原只sleep重试同死socket)
- 登录超时: SIGALRM给bs.login套30s闹钟,防永久挂起
- 断路器: 连续30只failed→save_progress+break+exit 2,不再假完成空跑2000票
- marker-resume: load_progress扫.baostock marker文件(非done JSON),failed无marker永远重试
- watcher: 每30min探测baostock,可登就续跑,冷却就等,exit 0则退出
9/9单测绿(Mac venv311 mock baostock)。
This commit is contained in:
2026-07-13 11:56:23 +08:00
parent b322f4e07b
commit 4b83452294
2 changed files with 197 additions and 23 deletions
+129 -23
View File
@@ -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 +1skipped 中性
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()
@@ -0,0 +1,68 @@
"""5yr 15min 回填冷却后自动续跑 watcher。
baostock 被限流/冷却时登录会卡死或失败。本 watcher 每 30min 探测一次 baostock
可登就启动 backfill_15min_baostockmarker-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()