699 lines
29 KiB
Python
699 lines
29 KiB
Python
#!/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:
|
||
# 2026-08-20 实锤: bs.login() 打印 "login success!" 后仍有后续往返, 服务端
|
||
# hiccup 时主线程可卡死在 login 内部 recv(08-20 19:26 挂死 50min+ 形态:
|
||
# logout✓ login success!打印后零输出) —— 与 fetch 同罩 _with_timeout。
|
||
lg = _with_timeout(bs.login, timeout=60)
|
||
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:
|
||
_with_timeout(bs.logout, timeout=30)
|
||
except Exception:
|
||
pass
|
||
if login_once():
|
||
return True
|
||
time.sleep(2)
|
||
try:
|
||
_with_timeout(bs.logout, timeout=30)
|
||
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 _valuation_frame(code, prefix, rows, fields):
|
||
"""估值 K 线行 -> 10 列 vdf [symbol,exchange,date,peTTM,psTTM,pcfNcfTTM,
|
||
pbMRQ,turn,pctChg,isST]。
|
||
|
||
fields 为请求字段串(列名对齐即可, 顺序无关): upsert_daily 传 DAILY_FIELDS
|
||
(15 列含 OHLCV), 回补传 VAL_BACKFILL_FIELDS (9 列纯估值)。
|
||
空串 peTTM/pbMRQ -> NaN (亏损/净资负), isST 字符串 -> int。
|
||
"""
|
||
df = pd.DataFrame(rows, columns=fields.split(","))
|
||
for c in ["turn", "pctChg", "peTTM", "psTTM", "pcfNcfTTM", "pbMRQ"]:
|
||
if c in df.columns:
|
||
df[c] = pd.to_numeric(df[c], errors="coerce")
|
||
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_MAP[prefix])
|
||
return vdf
|
||
|
||
|
||
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 = _valuation_frame(code, prefix, rows, DAILY_FIELDS)
|
||
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
|
||
|
||
|
||
# ======================== valuation_baostock 历史窗回补 (spec §19.11-a) ========================
|
||
# 2026.parquet 仅 08-13 起(bs 日喂起点), H1 缺口靠基线层 em valuation 异源口径补。
|
||
# 回补=按 (股, 窗口) 只拉估值 9 列 -> 中间产物落 data/_backfill_valuation/{tag}/
|
||
# (独立目录, 不进 sync_valuation_daily 的 scp -r 整目录拉取面) -> 终局与年文件
|
||
# merge 原子写回。断点序=buffer 先落盘、done 后写: 崩在中间最多整批重拉(键唯一幂等), 行永不丢。
|
||
|
||
VAL_BACKFILL_FIELDS = ("date,code,turn,pctChg,peTTM,psTTM,"
|
||
"pcfNcfTTM,pbMRQ,isST") # 无 OHLCV: 不动 dbbardata
|
||
BACKFILL_BATCH = 500 # 只/批: buffer 落盘+done 写入粒度
|
||
|
||
|
||
def merge_valuation_frames(old, new):
|
||
"""年文件 ∪ 回补段: drop_duplicates(symbol,date,keep='last')(与 upsert_daily
|
||
同语义) + sort。old=None(年文件不存在)时直接排序返回 new 副本。"""
|
||
if old is None:
|
||
return new.sort_values(["symbol", "date"]).reset_index(drop=True)
|
||
out = (pd.concat([old, new], ignore_index=True)
|
||
.drop_duplicates(["symbol", "date"], keep="last")
|
||
.sort_values(["symbol", "date"])
|
||
.reset_index(drop=True))
|
||
return out
|
||
|
||
|
||
def build_backfill_stock_list(current, snapshot):
|
||
"""当前全列表(全A含退市) ∪ 时点快照(query_all_stock) 并集去重保序。
|
||
|
||
时点兜底防 query_stock_basic 退市覆盖缺口; snapshot 空/None 只用当前列表。
|
||
"""
|
||
seen = set()
|
||
out = []
|
||
for item in list(current) + list(snapshot or []):
|
||
if item not in seen:
|
||
seen.add(item)
|
||
out.append(item)
|
||
return out
|
||
|
||
|
||
def midpoint_date(start, end):
|
||
"""窗口中点日(str YYYY-MM-DD) — query_all_stock 时点采样日。"""
|
||
s = dt.datetime.strptime(start, "%Y-%m-%d").date()
|
||
e = dt.datetime.strptime(end, "%Y-%m-%d").date()
|
||
return (s + (e - s) // 2).strftime("%Y-%m-%d")
|
||
|
||
|
||
def atomic_write_parquet(df, path):
|
||
"""tmp 写入 + os.replace 原子替换(2-3h 回补成果不能毁于写一半)。"""
|
||
tmp = Path(str(path) + ".tmp")
|
||
df.to_parquet(tmp, index=False)
|
||
os.replace(tmp, path)
|
||
|
||
|
||
def _fetch_all_stock_at(day):
|
||
"""query_all_stock(day) -> [(code, 'sh'/'sz')]; 非交易日/无数据返 []。"""
|
||
global QUERY_COUNT
|
||
QUERY_COUNT += 1
|
||
rs = bs.query_all_stock(day=day)
|
||
if rs.error_code != "0":
|
||
raise RuntimeError(f"query_all_stock: {rs.error_code} {rs.error_msg}")
|
||
idx = {n: i for i, n in enumerate(list(rs.fields))}
|
||
out = []
|
||
while rs.next():
|
||
r = rs.get_row_data()
|
||
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 _flush_buffer(buffer_path, frames):
|
||
"""批内 frames 并入 buffer.parquet(键唯一, 重拉幂等)。"""
|
||
new = (pd.concat(frames, ignore_index=True) if len(frames) > 1
|
||
else frames[0].reset_index(drop=True))
|
||
if buffer_path.exists():
|
||
old = pd.read_parquet(buffer_path)
|
||
new = pd.concat([old, new], ignore_index=True)
|
||
new.drop_duplicates(["symbol", "date"], keep="last").sort_values(
|
||
["symbol", "date"]).to_parquet(buffer_path, index=False)
|
||
|
||
|
||
def run_valuation_backfill(start, end):
|
||
"""一次性回补主流程: 拉取(断点续传) -> 终局 merge 原子写回。单年窗口。"""
|
||
global QUERY_COUNT
|
||
s = dt.datetime.strptime(start, "%Y-%m-%d").date()
|
||
e = dt.datetime.strptime(end, "%Y-%m-%d").date()
|
||
if s.year != e.year:
|
||
raise SystemExit(f"回补窗口须落在单一年份内: {start}~{end}")
|
||
tag = f"{s.year}{s.month:02d}{s.day:02d}_{e.year}{e.month:02d}{e.day:02d}"
|
||
work = VAL_DIR.parent / "_backfill_valuation" / tag
|
||
work.mkdir(parents=True, exist_ok=True)
|
||
done_path = work / "done.txt"
|
||
buffer_path = work / "buffer.parquet"
|
||
target = VAL_DIR / f"{e.year}.parquet"
|
||
log.info("[BACKFILL-VAL] window=%s~%s target=%s work=%s", start, end,
|
||
target, work)
|
||
|
||
if not login_with_retry():
|
||
log.error("[BACKFILL-VAL] 登录失败, exit 2")
|
||
sys.exit(2)
|
||
try:
|
||
try:
|
||
stocks = fetch_all_stocks_with_timeout()
|
||
except Exception as e_:
|
||
log.error("[BACKFILL-VAL] fetch_all: %s", e_)
|
||
sys.exit(1)
|
||
snap_day = midpoint_date(start, end)
|
||
try:
|
||
snapshot = _with_timeout(_fetch_all_stock_at, args=(snap_day,),
|
||
timeout=120)
|
||
except Exception as e_:
|
||
log.warning("query_all_stock(%s) 失败(跳过时点兜底): %s", snap_day, e_)
|
||
snapshot = []
|
||
stocks = build_backfill_stock_list(stocks, snapshot)
|
||
log.info("[BACKFILL-VAL] 名单=全A∪时点(%s): %d 只", snap_day, len(stocks))
|
||
|
||
done = set()
|
||
if done_path.exists():
|
||
done = {ln.strip() for ln in done_path.read_text(
|
||
encoding="utf-8").splitlines() if ln.strip()}
|
||
log.info("[BACKFILL-VAL] 断点续传: 已完成 %d 只", len(done))
|
||
pending = [(c, p) for c, p in stocks if c not in done]
|
||
|
||
stats = {"ok": 0, "empty": 0, "failed": 0, "rows": 0}
|
||
t0 = time.time()
|
||
batch_vdfs = [] # 批内行缓冲
|
||
batch_codes = [] # 批内 code(待 buffer 落盘后写 done)
|
||
stopped = False
|
||
with open(done_path, "a", encoding="utf-8") as done_f:
|
||
for i, (code, prefix) in enumerate(pending):
|
||
if QUERY_COUNT >= DAILY_LIMIT:
|
||
log.warning("query %d 达防线 %d, 剩余转下次续跑",
|
||
QUERY_COUNT, DAILY_LIMIT)
|
||
stopped = True
|
||
break
|
||
bs_code = f"{prefix}.{code}"
|
||
try:
|
||
rows = fetch_k_with_timeout(bs_code, VAL_BACKFILL_FIELDS,
|
||
"d", start, end)
|
||
except Exception as e_:
|
||
stats["failed"] += 1
|
||
if stats["failed"] <= 5 or stats["failed"] % 100 == 0:
|
||
log.warning("%s fetch err: %s", code, e_)
|
||
if not relogin():
|
||
log.error("%s relogin 失败, 跳过", code)
|
||
time.sleep(BS_INTERVAL)
|
||
continue
|
||
if rows:
|
||
vdf = _valuation_frame(code, prefix, rows,
|
||
VAL_BACKFILL_FIELDS)
|
||
# 防御: 窗口越界行过滤(接口契约上不会, 双保险)
|
||
vdf = vdf[(vdf["date"] >= start) & (vdf["date"] <= end)]
|
||
if len(vdf):
|
||
batch_vdfs.append(vdf)
|
||
stats["rows"] += len(vdf)
|
||
stats["ok"] += 1
|
||
else:
|
||
stats["empty"] += 1
|
||
batch_codes.append(code)
|
||
if (i + 1) % 100 == 0:
|
||
log.info("[BACKFILL-VAL] 进度 %d/%d ok=%d empty=%d "
|
||
"failed=%d rows=%d q=%d (%.0fs)",
|
||
i + 1, len(pending), stats["ok"], stats["empty"],
|
||
stats["failed"], stats["rows"], QUERY_COUNT,
|
||
time.time() - t0)
|
||
if (i + 1) % RELOGIN_EVERY == 0 and not relogin():
|
||
log.warning("周期 relogin 失败, 继续跑(下次被动 relogin 兜底)")
|
||
if len(batch_codes) >= BACKFILL_BATCH:
|
||
_flush_buffer(buffer_path, batch_vdfs)
|
||
done_f.write("\n".join(batch_codes) + "\n")
|
||
done_f.flush()
|
||
batch_vdfs, batch_codes = [], []
|
||
if i < len(pending) - 1:
|
||
time.sleep(BS_INTERVAL)
|
||
# 尾批: buffer 先落盘、done 后写(崩在中间=重拉幂等, 行永不丢)
|
||
if batch_codes:
|
||
if batch_vdfs:
|
||
_flush_buffer(buffer_path, batch_vdfs)
|
||
done_f.write("\n".join(batch_codes) + "\n")
|
||
log.info("[BACKFILL-VAL] 拉取完成 ok=%d empty=%d failed=%d rows=%d "
|
||
"query=%d 耗时%.0fs", stats["ok"], stats["empty"],
|
||
stats["failed"], stats["rows"], QUERY_COUNT, time.time() - t0)
|
||
|
||
# 终局 merge: buffer ∪ 年文件 -> 原子写回
|
||
if buffer_path.exists():
|
||
buf = pd.read_parquet(buffer_path)
|
||
old = pd.read_parquet(target) if target.exists() else None
|
||
out = merge_valuation_frames(old, buf)
|
||
atomic_write_parquet(out, target)
|
||
log.info("[BACKFILL-VAL] merge 完成 %s rows=%d (原 %s) 键唯一=%s",
|
||
target, len(out), len(old) if old is not None else 0,
|
||
out.duplicated(["symbol", "date"]).sum() == 0)
|
||
else:
|
||
log.warning("[BACKFILL-VAL] buffer 无数据(全空窗/全失败?), 不动 %s", target)
|
||
if stopped:
|
||
sys.exit(3)
|
||
finally:
|
||
try:
|
||
bs.logout()
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def _parse_args(argv=None):
|
||
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 重灌快)")
|
||
ap.add_argument("--backfill-valuation", nargs=2,
|
||
metavar=("START", "END"), default=None,
|
||
help="一次性: valuation_baostock 历史窗回补 "
|
||
"YYYY-MM-DD YYYY-MM-DD(只拉估值列, 不动 dbbardata)")
|
||
return ap.parse_args(argv)
|
||
|
||
|
||
def main():
|
||
global QUERY_COUNT
|
||
args = _parse_args()
|
||
if args.backfill_valuation:
|
||
start, end = args.backfill_valuation
|
||
run_valuation_backfill(start, end)
|
||
return
|
||
|
||
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()
|