diff --git a/scripts/data_platform/bs_5m_eod.py b/scripts/data_platform/bs_5m_eod.py index 76cac89..fc45302 100644 --- a/scripts/data_platform/bs_5m_eod.py +++ b/scripts/data_platform/bs_5m_eod.py @@ -17,18 +17,21 @@ - 库: BS_5M_DB 缺省 NAS 副本主库(VPS 同步目标, 也是 NAS 回测 provider 读的库, portfolio_worker.py 传的就是它); ensure_schema 已有唯一索引(NAS 侧名 uq_dbbardata) 则不重复建 -- 模式: 缺省每日增量(LOOKBACK 7 天, DSM 任务计划 ~19:05); --full 一次性回灌 - 2020-01-03+(baostock 分钟数据固定起点非滚动; 全区间每股恰好 1 次调用, 一趟 - ~5.5k query / ~3.6 亿行, 数小时级; per-stock 短事务, 中断重跑幂等) -- 配额: NAS 出口 IP 独立核算 48k/天(与 VPS 各自计数互不相干); 本脚本一趟远低 - 于限, DAILY_LIMIT 守卫保留 -- 单实例锁 bs_5m.lock(pid 活性检测): 防 DSM 定时与 --full 长跑撞车 +- 模式(2026-08-20 用户拍板"每天固定下一些", 拒绝一趟灌满): 每日一跑(DSM ~19:05)= + ①回灌未完成→下一片半年窗口(newest-first, 近端先到先可用; state 推进, 片没扫完 + 不推进次日重拉同片幂等) ②每日增量(LOOKBACK 7 天, 回灌期也跑, 近端始终新鲜)。 + 全程 2020-01-03(baostock 分钟固定起点)~今 ≈ 14 片 ≈ 两周补完, 每片 ~5.5k query / + ~3200 万行 / 3-4h; per-stock 短事务, 中断重跑幂等 +- 配额: NAS 出口 IP 独立核算 48k/天(与 VPS 各自计数互不相干); 每日 片+增量 + ~11k query, DAILY_LIMIT 守卫保留 +- 单实例锁 bs_5m.lock(pid 活性检测): 防手动补跑与 DSM 定时撞车 (同 IP 双 baostock 连接红线) 退出码: 0=完成/让路; 1=致命; 2=登录失败; 3=query 超限 graceful stop """ import argparse import datetime as dt +import json import logging import os import sqlite3 @@ -47,7 +50,7 @@ _DEFAULT_DB = "/volume1/stock/sanguo_vnpy_v2/data_backup/quant_trading.db" # NA DB = Path(os.environ.get("BS_5M_DB", _DEFAULT_DB)) LOOKBACK = int(os.environ.get("LOOKBACK_DAYS", "7")) FULL_START = "2020-01-03" # baostock 分钟数据固定起点(官网"近5年"口径 2020-01-03) -FULL_TIMEOUT = 300 # 全区间单股 ~6.4 万行, 默认 60s 超时不够 +CHUNK_DAYS = 183 # 回灌片宽(半年); 每片 ~5.5k query / ~3200 万行 / 3-4h M5_FIELDS = "date,time,code,open,high,low,close,volume,amount" logging.basicConfig(level=logging.INFO, @@ -56,13 +59,40 @@ logging.basicConfig(level=logging.INFO, log = logging.getLogger(__name__) -def calc_window(today, full): - """--full: 2020-01-03 固定起点全区间; 缺省: LOOKBACK 天增量窗口。""" - end = today.strftime("%Y-%m-%d") - if full: - return FULL_START, end - start = (today - dt.timedelta(days=LOOKBACK)).strftime("%Y-%m-%d") - return start, end +def _state_path(): + return DB.parent / "bs_5m_state.json" + + +def load_chunk_state(): + """{"next_end": "YYYY-MM-DD" | None}; 缺省 next_end=今天(首片=最近半年)。""" + if _state_path().exists(): + try: + return json.loads(_state_path().read_text(encoding="utf-8")) + except Exception: + log.warning("state 文件损坏, 回灌从最近半年重头(OR REPLACE 幂等无伤)") + return {"next_end": dt.date.today().isoformat()} + + +def save_chunk_state(state): + DB.parent.mkdir(parents=True, exist_ok=True) + _state_path().write_text(json.dumps(state, ensure_ascii=False, indent=2), + encoding="utf-8") + + +def next_chunk(today): + """下一个待回灌窗口 (start, end), newest-first 逐片向 2020 走; None=回灌完成。 + + 片宽 CHUNK_DAYS, 尾片与 FULL_START 对齐(前段有少量重叠, OR REPLACE 去重)。 + """ + state = load_chunk_state() + if state.get("next_end") is None: + return None + end = min(dt.date.fromisoformat(state["next_end"]), today) + if end <= dt.date.fromisoformat(FULL_START): + return None + start = max(dt.date.fromisoformat(FULL_START), + end - dt.timedelta(days=CHUNK_DAYS)) + return start.isoformat(), end.isoformat() def _has_unique_index(conn): @@ -166,8 +196,6 @@ def _release_lock(): def main(): ap = argparse.ArgumentParser() ap.add_argument("--limit", type=int, default=0, help="限股数(冒烟用)") - ap.add_argument("--full", action="store_true", - help="一次性回灌 2020-01-03+(缺省每日增量)") args = ap.parse_args() if not _acquire_lock(): @@ -179,11 +207,53 @@ def main(): _release_lock() +def _sweep(conn, stocks, start, end, swept, label): + """扫全 A 一个窗口: per-stock 短事务 + 周期 relogin + 间隔。 + + swept 为跨窗口累计股数(=query 数, 每股恰 1 次调用), 由调用方持有并回填; + 返 (stats, limit_reached)。 + """ + stats = {"ok": 0, "empty": 0, "failed": 0, "db_rows": 0} + limit_reached = False + t0 = time.time() + for i, (code, prefix) in enumerate(stocks): + if swept[0] >= DAILY_LIMIT: + log.warning("[%s] query %d 达防线 %d, graceful stop", + label, swept[0], DAILY_LIMIT) + limit_reached = True + break + swept[0] += 1 + try: + n = _process_one(conn, code, prefix, start, end, 60) + stats["db_rows"] += n + if n: + stats["ok"] += 1 + else: + stats["empty"] += 1 + except Exception as e: + stats["failed"] += 1 + if stats["failed"] <= 5 or stats["failed"] % 100 == 0: + log.warning("[%s] %s err: %s", label, code, e) + if not relogin(): + log.error("[%s] %s relogin 失败, 跳过", label, code) + if (i + 1) % RELOGIN_EVERY == 0: + log.info("[%s] 进度 %d/%d ok=%d empty=%d failed=%d rows=%d (%.0fs)", + label, i + 1, len(stocks), stats["ok"], stats["empty"], + stats["failed"], stats["db_rows"], time.time() - t0) + if not relogin(): + log.warning("[%s] 周期 relogin 失败, 继续跑", label) + if i < len(stocks) - 1: + time.sleep(BS_INTERVAL) + return stats, limit_reached + + def _run(args): today = dt.date.today() - start, end = calc_window(today, args.full) - log.info("bs_5m start window=%s~%s mode=%s db=%s", start, end, - "full" if args.full else "daily", DB) + chunk = next_chunk(today) + inc_start = (today - dt.timedelta(days=LOOKBACK)).strftime("%Y-%m-%d") + inc_end = today.strftime("%Y-%m-%d") + log.info("bs_5m start db=%s chunk=%s inc=%s~%s", DB, + "%s~%s" % chunk if chunk else "无(回灌完成)", inc_start, inc_end) if not login_with_retry(): log.error("[SKIP] 登录失败, exit 2") @@ -202,37 +272,28 @@ def _run(args): conn.execute("PRAGMA journal_mode = WAL") ensure_schema(conn) - timeout = FULL_TIMEOUT if args.full else 60 - stats = {"ok": 0, "empty": 0, "failed": 0, "db_rows": 0} + swept = [0] # 跨窗口累计 query 数(每股恰 1 次调用) limit_reached = False t0 = time.time() try: - for i, (code, prefix) in enumerate(stocks): - if i >= DAILY_LIMIT: # 每股恰好 1 query, i 即当日 query 计数 - log.warning("query %d 达防线 %d, graceful stop", i, DAILY_LIMIT) - limit_reached = True - break - try: - n = _process_one(conn, code, prefix, start, end, timeout) - stats["db_rows"] += n - if n: - stats["ok"] += 1 - else: - stats["empty"] += 1 - except Exception as e: - stats["failed"] += 1 - if stats["failed"] <= 5 or stats["failed"] % 100 == 0: - log.warning("%s err: %s", code, e) - if not relogin(): - log.error("%s relogin 失败, 跳过", code) - if (i + 1) % RELOGIN_EVERY == 0: - log.info("进度 %d/%d ok=%d empty=%d failed=%d rows=%d (%.0fs)", - i + 1, len(stocks), stats["ok"], stats["empty"], - stats["failed"], stats["db_rows"], time.time() - t0) - if not relogin(): - log.warning("周期 relogin 失败, 继续跑") - if i < len(stocks) - 1: - time.sleep(BS_INTERVAL) + if chunk: + cs, lr = _sweep(conn, stocks, chunk[0], chunk[1], swept, + "chunk %s~%s" % chunk) + limit_reached = limit_reached or lr + log.info("[CHUNK] %s~%s ok=%d empty=%d failed=%d rows=%d", + chunk[0], chunk[1], cs["ok"], cs["empty"], cs["failed"], + cs["db_rows"]) + if not lr and swept[0] < DAILY_LIMIT: + # 整片扫完才推进 state; 没扫完次日重拉同片(OR REPLACE 幂等续跑) + save_chunk_state({"next_end": chunk[0]}) + # 每日增量照跑(回灌期也跑, 近端 7 天始终新鲜) + istats, lr = _sweep(conn, stocks, inc_start, inc_end, swept, "inc") + limit_reached = limit_reached or lr + log.info("[INC] ok=%d empty=%d failed=%d rows=%d", + istats["ok"], istats["empty"], istats["failed"], + istats["db_rows"]) + if next_chunk(today) is None and not limit_reached: + save_chunk_state({"next_end": None}) # 回灌完结盖章(幂等) finally: conn.close() try: @@ -241,9 +302,7 @@ def _run(args): except Exception: pass - log.info("[DONE] mode=%s ok=%d empty=%d failed=%d db_rows=%d 耗时%.0fs", - "full" if args.full else "daily", stats["ok"], stats["empty"], - stats["failed"], stats["db_rows"], time.time() - t0) + log.info("[DONE] query=%d 耗时%.0fs", swept[0], time.time() - t0) sys.exit(3 if limit_reached else 0) diff --git a/tests/data_platform/test_bs_5m_eod.py b/tests/data_platform/test_bs_5m_eod.py index 7e100d1..20f25a3 100644 --- a/tests/data_platform/test_bs_5m_eod.py +++ b/tests/data_platform/test_bs_5m_eod.py @@ -10,6 +10,7 @@ INSERT OR IGNORE 去重(merge_increment.py), NAS 回测 reader 零改动可见 5 --full 全区间也一样); VPS 侧 bs_eod/bs_fund 用 VPS 自己的 IP 配额, 互不相干。 """ import datetime as dt +import json import os import sqlite3 import sys @@ -138,13 +139,32 @@ def test_process_one_success_writes(tmp_db_conn): assert n == 1 -# ---------- 窗口 ---------- +# ---------- 回灌分片(newest-first, 每天固定一片) ---------- -def test_calc_window_full_vs_daily(): +def test_next_chunk_first_is_recent_half_year(tmp_path, monkeypatch): + """无 state → 首片 = 最近半年(today-183, today], 近端先到先可用。""" + monkeypatch.setattr(b5, "DB", tmp_path / "dbbardata_5m.db") today = dt.date(2026, 8, 20) - assert b5.calc_window(today, full=True) == ("2020-01-03", "2026-08-20") - s, e = b5.calc_window(today, full=False) - assert s == "2026-08-13" and e == "2026-08-20" # LOOKBACK=7 + assert b5.next_chunk(today) == ("2026-02-18", "2026-08-20") + + +def test_next_chunk_walks_backward_and_finishes(tmp_path, monkeypatch): + """state 逐片向 2020 走; 尾片与 FULL_START 对齐; 越过即 None(回灌完成)。""" + monkeypatch.setattr(b5, "DB", tmp_path / "dbbardata_5m.db") + b5.save_chunk_state({"next_end": "2020-05-01"}) + assert b5.next_chunk(dt.date(2026, 8, 20)) == ("2020-01-03", "2020-05-01") + b5.save_chunk_state({"next_end": "2019-12-31"}) # 已越过起点 + assert b5.next_chunk(dt.date(2026, 8, 20)) is None + b5.save_chunk_state({"next_end": None}) # 完结盖章 + assert b5.next_chunk(dt.date(2026, 8, 20)) is None + + +def test_chunk_state_corrupt_falls_back_to_recent(tmp_path, monkeypatch): + """state 损坏 → 从最近半年重头(OR REPLACE 幂等无伤)。""" + monkeypatch.setattr(b5, "DB", tmp_path / "dbbardata_5m.db") + (tmp_path / "bs_5m_state.json").write_text("{broken", encoding="utf-8") + st = b5.load_chunk_state() + assert st["next_end"] == dt.date.today().isoformat() # ---------- 单实例锁(防 DSM 定时与全量回灌撞车) ---------- @@ -182,7 +202,11 @@ def test_main_runs_and_exits_0(main_env): with pytest.raises(SystemExit) as e: b5.main() assert e.value.code == 0 - assert m_proc.call_count == 2 + # 双 sweep: 回灌片(首跑=最近半年) + 每日增量, 各 2 股 → 4 次调用 + assert m_proc.call_count == 4 + # 整片扫完 → state 推进到片起点 + st = json.loads((main_env / "bs_5m_state.json").read_text(encoding="utf-8")) + assert st["next_end"] == (dt.date.today() - dt.timedelta(days=b5.CHUNK_DAYS)).isoformat() # 数据库文件+schema 已就位 conn = sqlite3.connect(str(b5.DB)) assert conn.execute( @@ -194,6 +218,16 @@ def test_main_runs_and_exits_0(main_env): assert not (main_env / "bs_5m.lock").exists() +def test_main_increment_only_when_backfill_done(main_env, monkeypatch): + """回灌完结盖章(next_end=None) → 只跑每日增量(单 sweep)。""" + b5.save_chunk_state({"next_end": None}) + with patch.object(b5, "_process_one", return_value=48) as m_proc: + with pytest.raises(SystemExit) as e: + b5.main() + assert e.value.code == 0 + assert m_proc.call_count == 2 # 仅增量 2 股 + + def test_main_exit2_when_login_fails(main_env, monkeypatch): monkeypatch.setattr(b5, "login_with_retry", MagicMock(return_value=False)) with pytest.raises(SystemExit) as e: