Files
sanguo_vnpy_v2/scripts/data_platform/bs_eod.py
T

463 lines
19 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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/<year>.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()