Files
claude_dev f94a145587 fix(data): bs_eod 卡死根治 — per-stock commit + baostock 超时包装 + 周期 relogin
根因(py-spy dump + netstat CLOSE_WAIT 实证): baostock 服务端关长连接→CLOSE_WAIT, send_msg 静默阻塞不抛异常, socket.setdefaulttimeout 不被 baostock 自己 socket 遵守, relogin 只在 error_code≠0 救不了; 一把大事务全程持 WAL 锁阻断全库。修复: per-stock commit 去大事务 + _with_timeout 线程超时包 fetch_k 打破静默 hang + 周期 relogin 每500主动刷连接。VPS --limit 3 验证 11s 不 hang。
2026-07-28 20:36:11 +08:00

354 lines
14 KiB
Python
Raw Permalink 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')
- 日线 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"
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()
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 _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_once():
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)
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()