diff --git a/scripts/data_platform/bs_eod.py b/scripts/data_platform/bs_eod.py index 0defa7c..26098fd 100644 --- a/scripts/data_platform/bs_eod.py +++ b/scripts/data_platform/bs_eod.py @@ -5,6 +5,7 @@ 每日收盘后跑(VPS, baostock 日终更新就绪): - 个股日线(含退市, LOOKBACK 7 天) -> dbbardata('d') INSERT OR REPLACE (治幸存者偏差) - 个股 15min(LOOKBACK 7) -> dbbardata('15m') +- 指数日K(2006+ 全量 REPLACE) -> dbbardata('d') 双源冗余, 主循环后跑不阻塞个股 - 日线 pe/pb/turn/pctChg/isST -> data/valuation_baostock/.parquet 追加 - DAILY_LIMIT=48000 单进程单登录, sleep 0.3s, login 探针 graceful skip @@ -50,6 +51,29 @@ DAILY_FIELDS = ("date,code,open,high,low,close,volume,amount,turn," "pctChg,peTTM,psTTM,pcfNcfTTM,pbMRQ,isST") M15_FIELDS = "date,time,code,open,high,low,close,volume,amount" +# ======================== 指数日K双源 (2026-08-20 P0) ======================== +# baostock 指数 2006+ 与 sina idx-eod 互为冗余: 同 PK(symbol,exchange,datetime, +# interval) INSERT OR REPLACE, 后写者胜 —— 治 000300 单源停更史(7-16)/932000 无点位 +# (baostock 也无 → 基准用 399303 国证2000 替代)/000938 停 2023(两源皆弃)。 +# 探针实证(2026-08-20 VPS): 下表代码全部有数; 932000/000938/000985/929/930/936/937 +# baostock 无 → 不进列表(免每日 warning 刷屏)。000016/399001/399006 未探针但属 +# 规模/成指类大概率有, rows=0 自动跳过不报错。 +INDEX_START = "2006-01-01" +INDEX_FIELDS = "date,code,open,high,low,close,volume,amount" +INDEX_CODES = ( + "sh.000001", # 上证综指 + "sh.000016", # 上证50 + "sh.000300", # 沪深300 + "sh.000905", # 中证500 + "sh.000852", # 中证1000 + "sh.000903", # 中证100 + "sz.399001", # 深证成指 + "sz.399006", # 创业板指 + "sz.399303", # 国证2000 (932000 中证2000 baostock 无, 以此作小盘基准) + # 中证一级行业(baostock 可用的 6 只: 928/931/932/933/934/935) + "sh.000928", "sh.000931", "sh.000932", "sh.000933", "sh.000934", "sh.000935", +) + logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s", handlers=[logging.StreamHandler(sys.stdout)]) log = logging.getLogger(__name__) @@ -267,6 +291,62 @@ def upsert_15m(conn, code, prefix, rows): return len(db) +def upsert_index_daily(conn, prefix, code, rows): + """指数日K rows -> dbbardata('d')。与 sina idx-eod 存量行同 PK, REPLACE 后写者胜。""" + if not rows: + return 0 + df = pd.DataFrame(rows, columns=INDEX_FIELDS.split(",")) + for c in ["open", "high", "low", "close", "volume", "amount"]: + df[c] = pd.to_numeric(df[c], errors="coerce") + exc = EXC_MAP[prefix] + db = pd.DataFrame({ + "symbol": code, "exchange": exc, + "datetime": df["date"].astype(str).map(normalize_daily_dt), + "interval": "d", "volume": df["volume"], "turnover": df["amount"], + "open_interest": 0.0, + "open_price": df["open"], "high_price": df["high"], + "low_price": df["low"], "close_price": df["close"], + }) + conn.executemany( + "INSERT OR REPLACE INTO dbbardata " + "(symbol,exchange,datetime,interval,volume,turnover,open_interest," + "open_price,high_price,low_price,close_price) VALUES (?,?,?,?,?,?,?,?,?,?,?)", + db.itertuples(index=False, name=None)) + return len(db) + + +def run_index_eod(conn, end): + """主循环后拉指数日K(辅助数据层): 单指数 fetch/upsert 失败只 log 不外抛。 + + 指数是辅助数据 —— 任何异常不得影响个股 EOD 的退出码(schtask 结果码语义保持)。 + fetch_k_with_timeout 内部已计 QUERY_COUNT(指数 ~15 query/天, 预算可忽略)。 + """ + n_ok = n_empty = 0 + for bs_code in INDEX_CODES: + prefix, code = bs_code.split(".", 1) + try: + rows = fetch_k_with_timeout(bs_code, INDEX_FIELDS, "d", INDEX_START, end) + except Exception as e: + log.warning("指数 %s fetch err: %s", bs_code, e) + if not relogin(): + log.error("指数段 relogin 失败, 提前结束本段") + return + continue + try: + with conn: + n = upsert_index_daily(conn, prefix, code, rows) + except Exception as e: + log.warning("指数 %s upsert err: %s", bs_code, e) + continue + if n: + n_ok += 1 + else: + n_empty += 1 + log.warning("指数 %s 返回 0 行(源缺该指数?), 跳过", bs_code) + log.info("[INDEX] ok=%d empty=%d (baostock 与 sina idx-eod 双源互备)", + n_ok, n_empty) + + def _process_one_stock(conn, code, prefix, args, start, end): """单只股票: fetch_k + upsert, 在 with conn 短事务里执行 (大事务根治). @@ -357,6 +437,14 @@ def main(): log.warning("周期 relogin 失败, 继续跑 (下次 fetch 失败时被动 relogin 兜底)") if i < len(stocks) - 1: time.sleep(BS_INTERVAL) + # 指数日K双源段(主循环后): 达限跳过守预算; 段内异常全吞, 不改退出码 + if not limit_reached: + try: + run_index_eod(conn, end) + except Exception as e: + log.error("[INDEX] 段级异常(不影响个股 EOD 结果): %s", e) + else: + log.warning("query 达限, 跳过指数段(次日 2006+ 全量 REPLACE 自愈)") finally: conn.close() try: diff --git a/tests/data_platform/test_bs_eod_index.py b/tests/data_platform/test_bs_eod_index.py new file mode 100644 index 0000000..5b5a8eb --- /dev/null +++ b/tests/data_platform/test_bs_eod_index.py @@ -0,0 +1,203 @@ +# -*- coding: utf-8 -*- +"""TDD for bs_eod.py 指数日K双源段 (2026-08-20 P0 数据补全). + +背景: 指数点位此前单源 sina idx-eod, 有停更史(000300 曾停 7-16, 000938 停 2023, +932000 三源无). 本段在个股 EOD 主循环后追加 baostock 指数日K(2006+ 全量 REPLACE), +与 sina 互为冗余 —— 同 PK(symbol,exchange,datetime,interval) 后写者胜, 唯一索引防双行. + +探针实证(2026-08-20 VPS): INDEX_CODES 所列全部有数; 932000/000938/000985/ +929/930/936/937 baostock 无 → 不进列表. 000016/399001/399006 未探针, rows=0 跳过. + +测试覆盖: + 1. upsert_index_daily: sh/sz 映射 SSE/SZSE + datetime 归一 + REPLACE 不双行 + 2. run_index_eod: 空/异常容错(单指数失败不炸段), relogin 兜底 + 3. INDEX_CODES 卫生: 399303 在, 932000 不在, 格式合法 + 4. main() 集成: 正常跑完调指数段; query 达限跳过指数段(守 DAILY_LIMIT 预算) +""" +import re +import sqlite3 +import sys +from unittest.mock import MagicMock, patch + +import pytest + +if "baostock" not in sys.modules: + sys.modules["baostock"] = MagicMock() + +from scripts.data_platform import bs_eod # noqa: E402 + + +# ---------- Fixtures ---------- + +@pytest.fixture +def tmp_db(tmp_path): + """临时 sqlite DB 带 dbbardata 表 + 生产同款唯一索引 (REPLACE 去重前提).""" + db_path = tmp_path / "test.db" + conn = sqlite3.connect(str(db_path)) + conn.execute( + "CREATE TABLE dbbardata (" + "symbol TEXT, exchange TEXT, datetime TEXT, interval TEXT, " + "volume REAL, turnover REAL, open_interest REAL, " + "open_price REAL, high_price REAL, low_price REAL, close_price REAL)" + ) + conn.execute( + "CREATE UNIQUE INDEX dbbardata_symbol_exchange_interval_datetime " + "ON dbbardata (symbol, exchange, interval, datetime)" + ) + conn.commit() + yield conn + conn.close() + + +@pytest.fixture +def reset_query_count(): + original = bs_eod.QUERY_COUNT + bs_eod.QUERY_COUNT = 0 + yield + bs_eod.QUERY_COUNT = original + + +def _idx_rows(bs_code, date="2026-08-18", close="4725.813"): + """模拟 baostock 指数日K row (8 列 INDEX_FIELDS).""" + return [[date, bs_code, "4700.1", "4750.0", "4690.0", close, + "123456789", "520000000000"]] + + +# ---------- upsert_index_daily ---------- + +def test_upsert_index_maps_sh_to_sse(tmp_db): + """sh.000300 -> symbol=000300/exchange=SSE/interval=d, datetime 归一纯日期.""" + n = bs_eod.upsert_index_daily(tmp_db, "sh", "000300", _idx_rows("sh.000300")) + tmp_db.commit() + assert n == 1 + row = tmp_db.execute( + "SELECT symbol, exchange, datetime, interval, close_price FROM dbbardata" + ).fetchone() + assert row[0] == "000300" + assert row[1] == "SSE" + assert row[2] == "2026-08-18" # 纯日期, 与 sina idx-eod 存量行同格式 + assert row[3] == "d" + assert abs(row[4] - 4725.813) < 1e-6 + + +def test_upsert_index_maps_sz_to_szse(tmp_db): + """sz.399303 国证2000 -> exchange=SZSE (932000 无点位的替代基准).""" + bs_eod.upsert_index_daily(tmp_db, "sz", "399303", _idx_rows("sz.399303")) + tmp_db.commit() + row = tmp_db.execute("SELECT exchange FROM dbbardata").fetchone() + assert row[0] == "SZSE" + + +def test_upsert_index_replaces_same_pk_no_dup(tmp_db): + """同 PK 重复写 -> REPLACE 不双行, 后写者胜 (与 sina idx-eod 互备的根基).""" + bs_eod.upsert_index_daily(tmp_db, "sh", "000300", + _idx_rows("sh.000300", close="4725.813")) + bs_eod.upsert_index_daily(tmp_db, "sh", "000300", + _idx_rows("sh.000300", close="4600.0")) + tmp_db.commit() + rows = tmp_db.execute( + "SELECT close_price FROM dbbardata ORDER BY close_price" + ).fetchall() + assert len(rows) == 1 # 唯一索引 + REPLACE: 永不双行 + assert abs(rows[0][0] - 4600.0) < 1e-6 # 后写者胜 + + +def test_upsert_index_empty_rows_no_op(tmp_db): + """空 rows 返 0 不写 (源缺该指数的常态).""" + n = bs_eod.upsert_index_daily(tmp_db, "sh", "000300", []) + assert n == 0 + assert tmp_db.execute("SELECT COUNT(*) FROM dbbardata").fetchone()[0] == 0 + + +# ---------- run_index_eod 容错 ---------- + +def test_run_index_empty_result_skipped_others_continue(tmp_db, reset_query_count): + """某指数 0 行(源缺)只 warning, 其余继续 —— 单指数缺失不炸段.""" + def fake_fetch(bs_code, fields, freq, start, end): + assert freq == "d" and start == bs_eod.INDEX_START + return _idx_rows(bs_code) if bs_code == "sh.000300" else [] + with patch.object(bs_eod, "fetch_k_with_timeout", side_effect=fake_fetch): + bs_eod.run_index_eod(tmp_db, "2026-08-19") + rows = tmp_db.execute( + "SELECT symbol FROM dbbardata WHERE symbol='000300'" + ).fetchall() + assert len(rows) == 1 + + +def test_run_index_fetch_error_relogin_then_continue(tmp_db, reset_query_count): + """单指数 fetch 异常 -> relogin 兜底, 段不中断.""" + def fake_fetch(bs_code, fields, freq, start, end): + if bs_code == bs_eod.INDEX_CODES[0]: + raise RuntimeError("10002007 网络接收错误") + return _idx_rows(bs_code) + with patch.object(bs_eod, "fetch_k_with_timeout", side_effect=fake_fetch), \ + patch.object(bs_eod, "relogin", return_value=True) as m_rel: + bs_eod.run_index_eod(tmp_db, "2026-08-19") + assert m_rel.call_count == 1 + # 后续指数仍写入 + assert tmp_db.execute("SELECT COUNT(*) FROM dbbardata").fetchone()[0] >= 1 + + +def test_run_index_upsert_error_swallowed(tmp_db, reset_query_count): + """upsert 异常被吞(只 log), 段继续 —— 指数是辅助数据, 不拖垮个股 EOD 退出码.""" + with patch.object(bs_eod, "fetch_k_with_timeout", + side_effect=lambda *a, **k: _idx_rows("sh.000300")), \ + patch.object(bs_eod, "upsert_index_daily", + side_effect=RuntimeError("disk full")): + bs_eod.run_index_eod(tmp_db, "2026-08-19") # 不应 raise + + +# ---------- INDEX_CODES 卫生 ---------- + +def test_index_codes_probed_availability(): + """列表只含探针实证有数的: 399303 在(932000 替代), 932000/000938 不在.""" + assert "sz.399303" in bs_eod.INDEX_CODES + assert "sh.000300" in bs_eod.INDEX_CODES + for absent in ("sh.932000", "sh.000938", "sh.000985", + "sh.000929", "sh.000930", "sh.000936", "sh.000937"): + assert absent not in bs_eod.INDEX_CODES, f"{absent} baostock 无, 不应每日拉" + for code in bs_eod.INDEX_CODES: + assert re.match(r"^(sh|sz)\.\d{6}$", code), f"非法代码格式: {code}" + + +# ---------- main() 集成 ---------- + +@pytest.fixture +def isolated_main_env(tmp_path, monkeypatch, reset_query_count): + monkeypatch.setattr(bs_eod, "DB", tmp_path / "fake.db") + val = tmp_path / "val" + val.mkdir(parents=True, exist_ok=True) + monkeypatch.setattr(bs_eod, "VAL_DIR", val) + monkeypatch.setattr(sys, "argv", ["bs_eod.py"]) + mock_bs = MagicMock() + mock_bs.login.return_value.error_code = "0" + monkeypatch.setattr(bs_eod, "bs", mock_bs) + monkeypatch.setattr(bs_eod, "BS_INTERVAL", 0.0) + return tmp_path + + +def test_main_runs_index_section_after_stocks(isolated_main_env): + """正常跑完个股 -> 指数段被调用一次, 退出码 0.""" + stocks = [("600001", "sh")] + with patch.object(bs_eod, "fetch_all_stocks_with_timeout", + return_value=stocks), \ + patch.object(bs_eod, "_process_one_stock", return_value=(1, 0)), \ + patch.object(bs_eod, "run_index_eod") as m_idx: + with pytest.raises(SystemExit) as exc: + bs_eod.main() + assert exc.value.code == 0 + m_idx.assert_called_once() + + +def test_main_skips_index_when_query_limit_reached(isolated_main_env, monkeypatch): + """DAILY_LIMIT=0 -> 个股循环立即达限 graceful stop(exit 3), 指数段跳过守预算.""" + monkeypatch.setattr(bs_eod, "DAILY_LIMIT", 0) + stocks = [("600001", "sh")] + with patch.object(bs_eod, "fetch_all_stocks_with_timeout", + return_value=stocks), \ + patch.object(bs_eod, "_process_one_stock", return_value=(1, 0)), \ + patch.object(bs_eod, "run_index_eod") as m_idx: + with pytest.raises(SystemExit) as exc: + bs_eod.main() + assert exc.value.code == 3 + m_idx.assert_not_called() diff --git a/tests/data_platform/test_bs_eod_resilience.py b/tests/data_platform/test_bs_eod_resilience.py index 6c71c26..e2514c9 100644 --- a/tests/data_platform/test_bs_eod_resilience.py +++ b/tests/data_platform/test_bs_eod_resilience.py @@ -418,6 +418,7 @@ def test_periodic_relogin_called_every_n_stocks(isolated_main_env, monkeypatch): stocks = [(f"60000{i}", "sh") for i in range(6)] # 6 只 / 每 2 只 → 3 次 with patch.object(bs_eod, "fetch_all_stocks_with_timeout", return_value=stocks), \ patch.object(bs_eod, "_process_one_stock", return_value=(1, 0)), \ + patch.object(bs_eod, "run_index_eod"), \ patch.object(bs_eod, "relogin", return_value=True) as m_rel: with pytest.raises(SystemExit) as exc: bs_eod.main() @@ -432,6 +433,7 @@ def test_periodic_relogin_failure_does_not_crash_main(isolated_main_env, monkeyp stocks = [(f"60000{i}", "sh") for i in range(4)] with patch.object(bs_eod, "fetch_all_stocks_with_timeout", return_value=stocks), \ patch.object(bs_eod, "_process_one_stock", return_value=(1, 0)) as m_proc, \ + patch.object(bs_eod, "run_index_eod"), \ patch.object(bs_eod, "relogin", return_value=False) as m_rel: with pytest.raises(SystemExit) as exc: bs_eod.main() @@ -447,6 +449,7 @@ def test_periodic_relogin_disabled_when_relogin_every_huge(isolated_main_env, mo stocks = [(f"60000{i}", "sh") for i in range(5)] with patch.object(bs_eod, "fetch_all_stocks_with_timeout", return_value=stocks), \ patch.object(bs_eod, "_process_one_stock", return_value=(1, 0)), \ + patch.object(bs_eod, "run_index_eod"), \ patch.object(bs_eod, "relogin", return_value=True) as m_rel: with pytest.raises(SystemExit): bs_eod.main()