#!/usr/bin/env python3 # -*- coding: utf-8 -*- """bs_eod.py — sanguo-bs-eod (方案A schtask 18:05): baostock 个股 EOD 增量。 每日收盘后跑(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 预算: 5537股 × (1日线+1 15min) ≈ 11000 query/天 = 48000 的 23%, 安全。 退出码: 0=完成; 1=致命; 2=黑名单 graceful skip; 3=query 超限 graceful stop """ import argparse import datetime as dt import logging import os import socket import sys import threading import time from pathlib import Path for _k in ("http_proxy", "https_proxy", "HTTP_PROXY", "HTTPS_PROXY", "all_proxy", "ALL_PROXY"): os.environ.pop(_k, None) socket.setdefaulttimeout(30) try: sys.stdout.reconfigure(line_buffering=True) except (AttributeError, ValueError): pass import baostock as bs import pandas as pd from dbbardata_utils import normalize_daily_dt BASE = Path(r"C:\sanguo_vnpy_v2") DB = BASE / "data" / "quant_trading.db" VAL_DIR = BASE / "data" / "valuation_baostock" LOOKBACK = int(os.environ.get("LOOKBACK_DAYS", "7")) DAILY_LIMIT = int(os.environ.get("BS_DAILY_LIMIT", "48000")) # ③ 周期 relogin: 每 N 只主动 relogin, 防服务端长连接 idle 超时 → CLOSE_WAIT 静默 hang. # CLOSE_WAIT 是静默阻塞不抛异常, 被动 relogin (error_code != 0 触发) 救不了, 必须周期主动刷连接. RELOGIN_EVERY = int(os.environ.get("BS_RELOGIN_EVERY", "500")) BS_INTERVAL = 0.3 QUERY_COUNT = 0 EXC_MAP = {"sh": "SSE", "sz": "SZSE"} 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__) def login_once(): try: lg = bs.login() if lg.error_code == "0": return True log.error("login fail: %s %s", lg.error_code, lg.error_msg) return False except Exception as e: log.error("login exc: %s", e) return False def relogin(): try: bs.logout() except Exception: pass if login_once(): return True time.sleep(2) try: bs.logout() except Exception: pass return login_once() # 启动登录重试间隔(秒): 1 + 2 轮重试, 共 3 次尝试。 # 2026-08-20 用户拍板: 优先自身健壮性(实测无频繁重连/超 48k 配额问题, 封禁风险低)。 # 08-18 实录: 18:05:06 登录被 10054 掐断即 graceful exit, 整天缺口裸奔到次日; # 瞬时闪断秒~分钟级恢复, 2/5 分钟两轮串行重试(无并发)覆盖; 仍失败(真服务故障) # 才放弃 exit 2, 次日 LOOKBACK=7 自愈。 LOGIN_RETRY_DELAYS = (120, 300) def login_with_retry(): """启动登录: 失败按 LOGIN_RETRY_DELAYS 串行重试(治瞬时闪断), 全败返 False。""" if login_once(): return True total = len(LOGIN_RETRY_DELAYS) + 1 for i, delay in enumerate(LOGIN_RETRY_DELAYS, start=2): log.warning("登录失败, 第 %d/%d 次重试(等待 %ds)", i, total, delay) time.sleep(delay) if login_once(): return True return False def fetch_all_stocks(): """query_stock_basic() 无参 -> 全 A 含退市 (type=1)。返回 [(code, 'sh'/'sz')]。""" global QUERY_COUNT QUERY_COUNT += 1 rs = bs.query_stock_basic() if rs.error_code != "0": raise RuntimeError(f"query_stock_basic: {rs.error_code} {rs.error_msg}") fields = list(rs.fields) idx = {n: i for i, n in enumerate(fields)} out = [] while rs.next(): r = rs.get_row_data() if r[idx["type"]] != "1": continue bc = r[idx["code"]] if "." not in bc: continue prefix, num = bc.split(".", 1) if prefix in ("sh", "sz") and len(num) == 6 and num.isdigit(): out.append((num, prefix)) return out def fetch_k(bs_code, fields, freq, start, end): global QUERY_COUNT QUERY_COUNT += 1 rs = bs.query_history_k_data_plus(bs_code, fields, start_date=start, end_date=end, frequency=freq, adjustflag="3") if rs.error_code != "0": raise RuntimeError(f"{bs_code}: {rs.error_code} {rs.error_msg}") rows = [] while rs.next(): rows.append(rs.get_row_data()) return rows # ======================== 超时包装 ======================== # 根因: socket.setdefaulttimeout(30) 对 baostock 客户端的 rs.next() / recv 不可靠遵守, # 服务端 hiccup 会无限阻塞主循环. 用 daemon 工作线程 + join(timeout) 强制上限. # # 线程方案 vs subprocess 隔离: 选线程. # 理由: (1) baostock 模块级 singleton, 主线程 join 等子线程 → 单线程串行调用安全; # (2) 超时后子线程 daemon 化泄漏, 后续 relogin 的 logout 关旧 socket → 旧线程 # recv 报错自死, 新 login 走新 socket 不受污染 (relogin retry 兜底恢复); # (3) subprocess 方案需重新 login (~1s) + 序列化 rows 复杂, 不值; # (4) 与 akshare_static_download.call_ak_with_timeout 同范式. def _with_timeout(fn, args=(), kwargs=None, timeout=60): """daemon 线程跑 fn(*args, **kwargs), timeout 秒未完成 raise TimeoutError. 超时后工作线程泄漏 (daemon=True, 进程退出时强杀); 主线程立即返回让上层 relogin. 子线程异常透传给主线程 (BaseException 也捕获, 避免 daemon 吞 KeyboardInterrupt). """ if kwargs is None: kwargs = {} box = {"val": None, "exc": None} def worker(): try: box["val"] = fn(*args, **kwargs) except BaseException as e: # noqa: BLE001 - 透传所有异常含 KeyboardInterrupt box["exc"] = e t = threading.Thread(target=worker, daemon=True) t.start() t.join(timeout) if t.is_alive(): raise TimeoutError(f"{getattr(fn, '__name__', repr(fn))} 超过 {timeout}s") if box["exc"] is not None: raise box["exc"] return box["val"] def fetch_k_with_timeout(bs_code, fields, freq, start, end, timeout=60): """fetch_k + 超时保护 (默认 60s; baostock hiccup 不再无限阻塞).""" return _with_timeout( fetch_k, args=(bs_code, fields, freq, start, end), timeout=timeout, ) def fetch_all_stocks_with_timeout(timeout=120): """fetch_all_stocks + 超时保护 (全 A 列表一次性返回, 给 120s).""" return _with_timeout(fetch_all_stocks, timeout=timeout) def upsert_daily(conn, code, prefix, rows): """日线 rows -> dbbardata('d') + valuation_baostock 当年 parquet 追加。""" if not rows: return 0 df = pd.DataFrame(rows, columns=DAILY_FIELDS.split(",")) for c in ["open", "high", "low", "close", "volume", "amount", "turn", "pctChg", "peTTM", "psTTM", "pcfNcfTTM", "pbMRQ"]: df[c] = pd.to_numeric(df[c], errors="coerce") exc = EXC_MAP[prefix] # OHLCV -> dbbardata('d') — datetime 归一纯日期 (dbbardata 双行根治方案A) 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)) # pe/pb -> parquet 追加 (isST->int, 修 pyarrow ArrowTypeError) vdf = df[["date", "peTTM", "psTTM", "pcfNcfTTM", "pbMRQ", "turn", "pctChg", "isST"]].copy() vdf["isST"] = pd.to_numeric(vdf["isST"], errors="coerce").fillna(0).astype(int) vdf.insert(0, "symbol", code) vdf.insert(1, "exchange", exc) yr = dt.date.today().year p = VAL_DIR / f"{yr}.parquet" if p.exists(): try: old = pd.read_parquet(p) vdf = pd.concat([old, vdf]).drop_duplicates(["symbol", "date"], keep="last") except Exception: pass vdf.sort_values(["symbol", "date"]).to_parquet(p, index=False) return len(db) def _build_15m_dt(date_series, time_series): """baostock 15min datetime 拼接: date="2026-07-21"(带 -) + time="20260721094500000"(17 位)。 从 17 位 time 第 8-12 位提取 HHMM, date 直连(带 -)。 产出 'YYYY-MM-DD HH:MM:00' (符合 dbbardata GLOB 模式, 不被清理误删)。 bug 根因(2026-07-25): 原代码假设 date 纯数字 + time[:6] 取年月, 但 baostock 实测 date 带 -, time[:6]=YYYYMM, 产乱 datetime 致 15min 全市场停 7-17。 """ _t = time_series.astype(str) return (date_series.astype(str) + " " + _t.str.slice(8, 10) + ":" + _t.str.slice(10, 12) + ":00") def upsert_15m(conn, code, prefix, rows): if not rows: return 0 df = pd.DataFrame(rows, columns=M15_FIELDS.split(",")) for c in ["open", "high", "low", "close", "volume", "amount"]: df[c] = pd.to_numeric(df[c], errors="coerce") exc = EXC_MAP[prefix] dt_col = _build_15m_dt(df["date"], df["time"]) db = pd.DataFrame({ "symbol": code, "exchange": exc, "datetime": dt_col, "interval": "15m", "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 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 短事务里执行 (大事务根治). 每只股票一个事务 — hang/kill/异常最多丢 1 只, 已 commit 的其他股不受影响. fetch 也包在事务里 (task 要求: "fetch_k + upsert 包在自己事务里"); fetch 用 fetch_k_with_timeout 保护, 网络挂最多锁 timeout 秒. 成功返 (n1_daily, n2_15m); 异常时 with conn 自动 ROLLBACK 该股, 异常上抛. """ bs_code = f"{prefix}.{code}" n1 = 0 n2 = 0 with conn: # 显式短事务: 成功 commit / 异常 rollback (per-stock 原子) if not args.no_daily: d_rows = fetch_k_with_timeout(bs_code, DAILY_FIELDS, "d", start, end) n1 = upsert_daily(conn, code, prefix, d_rows) if not args.no_15m: m_rows = fetch_k_with_timeout(bs_code, M15_FIELDS, "15", start, end) n2 = upsert_15m(conn, code, prefix, m_rows) return n1, n2 def main(): global QUERY_COUNT ap = argparse.ArgumentParser() ap.add_argument("--limit", type=int, default=0) ap.add_argument("--no-15m", action="store_true") ap.add_argument("--no-daily", action="store_true", help="跳日线, 只跑 15min(用于 15min 重灌快)") args = ap.parse_args() today = dt.date.today() end = today.strftime("%Y-%m-%d") start = (today - dt.timedelta(days=LOOKBACK)).strftime("%Y-%m-%d") log.info("bs_eod start window=%s~%s LOOKBACK=%d limit=%s", start, end, LOOKBACK, args.limit or "无") if not login_with_retry(): log.error("[SKIP] 重试后仍登录失败(baostock 冷却/服务故障), graceful exit 2") sys.exit(2) try: stocks = fetch_all_stocks_with_timeout() except Exception as e: log.error("[FATAL] fetch_all: %s", e) sys.exit(1) log.info("全 A 含退市: %d 只", len(stocks)) if args.limit: stocks = stocks[:args.limit] import sqlite3 VAL_DIR.mkdir(parents=True, exist_ok=True) conn = sqlite3.connect(str(DB), timeout=60) conn.execute("PRAGMA busy_timeout = 60000") conn.execute("PRAGMA journal_mode = WAL") stats = {"ok": 0, "empty": 0, "failed": 0, "db_rows": 0} limit_reached = False t0 = time.time() # 大事务根治: 不再 BEGIN/COMMIT 包全程. 每只股票 with conn 短事务独立提交, # hang/kill/崩溃最多丢 1 只 (per-stock 隔离), 已 commit 的进度不丢. try: for i, (code, prefix) in enumerate(stocks): if QUERY_COUNT >= DAILY_LIMIT: log.warning("query %d 达防线 %d, graceful stop", QUERY_COUNT, DAILY_LIMIT) limit_reached = True break try: n1, n2 = _process_one_stock(conn, code, prefix, args, start, end) stats["db_rows"] += n1 + n2 if n1 + n2: stats["ok"] += 1 else: stats["empty"] += 1 except Exception as e: # _process_one_stock 异常已 rollback 该股, 其他股已 commit 不受影响 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 q=%d (%.0fs)", i + 1, len(stocks), stats["ok"], stats["empty"], stats["failed"], QUERY_COUNT, time.time() - t0) # ③ 主动 relogin (双保险之治本): 服务端长连接 idle 超时 → CLOSE_WAIT 静默 hang, # 被动 relogin 不触发 (不抛异常), 必须周期主动 logout+login 刷新连接. # ② _with_timeout 是治标兜底, 真挂了能打破; 这里治本避免走到那一步. if not relogin(): 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: bs.logout() except Exception: pass log.info("[DONE] ok=%d empty=%d failed=%d db_rows=%d query=%d 耗时%.0fs", stats["ok"], stats["empty"], stats["failed"], stats["db_rows"], QUERY_COUNT, time.time() - t0) sys.exit(3 if limit_reached else 0) if __name__ == "__main__": main()