Files
sanguo_vnpy_v2/scripts/data_platform/bs_5m_eod.py
T

311 lines
12 KiB
Python

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""bs_5m_eod.py — NAS 独立 5 分钟线备份 (数据补全 2026-08-20 用户拍板, 同日二次拍板改同库).
定位: 5m 直接写进 NAS 副本主库 dbbardata 表(与 VPS 同步镜像同一个 db/同一张表),
靠 interval='5m' 与镜像行(d/15m)在唯一键上正交实现互不干扰 —— 用户目标: NAS 就地
5m 回测/回放(reader 零改动), 未来 VPS 扩容直接把 interval='5m' 导出 merge 过去。
同库安全性(结构性, 非约定):
- 同步链(export_increment/merge_increment/sync_dbbardata)导出**不带 id 列**, NAS 侧
INSERT OR IGNORE 靠 UNIQUE(symbol,exchange,interval,datetime) 去重 → 5m 行与镜像行
键永不相交, 双向都碰不到对方; NAS 本地行不会被同步删除(纯 merge 无文件覆盖)
- id 由 NAS AUTOINCREMENT 自分配, 不存在跨库 id 冲突
- 写锁竞争: 双写者(WAL+busy_timeout)串行化; 与每日增量 merge(秒级)重叠时单股失败
次日 LOOKBACK=7 自愈; 全量 merge 分片(分特级)期间失败股同机制自愈
- 库: BS_5M_DB 缺省 NAS 副本主库(VPS 同步目标, 也是 NAS 回测 provider 读的库,
portfolio_worker.py 传的就是它); ensure_schema 已有唯一索引(NAS 侧名 uq_dbbardata)
则不重复建
- 模式(2026-08-20 用户拍板"每天固定下一些", 拒绝一趟灌满): 每日一跑(DSM ~19:05)=
①回灌未完成→下一片半年窗口(newest-first, 近端先到先可用; state 推进, 片没扫完
不推进次日重拉同片幂等) ②每日增量(LOOKBACK 7 天, 回灌期也跑, 近端始终新鲜)。
全程 2020-01-03(baostock 分钟固定起点)~今 ≈ 14 片 ≈ 两周补完, 每片 ~5.5k query /
~3200 万行 / 3-4h; per-stock 短事务, 中断重跑幂等
- 配额: NAS 出口 IP 独立核算 48k/天(与 VPS 各自计数互不相干); 每日 片+增量
~11k query, DAILY_LIMIT 守卫保留
- 单实例锁 bs_5m.lock(pid 活性检测): 防手动补跑与 DSM 定时撞车
(同 IP 双 baostock 连接红线)
退出码: 0=完成/让路; 1=致命; 2=登录失败; 3=query 超限 graceful stop
"""
import argparse
import datetime as dt
import json
import logging
import os
import sqlite3
import sys
import time
from pathlib import Path
import pandas as pd
from bs_eod import ( # noqa: E402 — 复用同一套健壮性封装, 口径与 bs_eod 15m 完全一致
BS_INTERVAL, DAILY_LIMIT, EXC_MAP, RELOGIN_EVERY, _build_15m_dt,
fetch_all_stocks_with_timeout, fetch_k_with_timeout, login_with_retry, relogin,
)
_DEFAULT_DB = "/volume1/stock/sanguo_vnpy_v2/data_backup/quant_trading.db" # NAS 副本主库(同步目标=回测源, 同库正交共存)
DB = Path(os.environ.get("BS_5M_DB", _DEFAULT_DB))
LOOKBACK = int(os.environ.get("LOOKBACK_DAYS", "7"))
FULL_START = "2020-01-03" # baostock 分钟数据固定起点(官网"近5年"口径 2020-01-03)
CHUNK_DAYS = 183 # 回灌片宽(半年); 每片 ~5.5k query / ~3200 万行 / 3-4h
M5_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 _state_path():
return DB.parent / "bs_5m_state.json"
def load_chunk_state():
"""{"next_end": "YYYY-MM-DD" | None}; 缺省 next_end=今天(首片=最近半年)。"""
if _state_path().exists():
try:
return json.loads(_state_path().read_text(encoding="utf-8"))
except Exception:
log.warning("state 文件损坏, 回灌从最近半年重头(OR REPLACE 幂等无伤)")
return {"next_end": dt.date.today().isoformat()}
def save_chunk_state(state):
DB.parent.mkdir(parents=True, exist_ok=True)
_state_path().write_text(json.dumps(state, ensure_ascii=False, indent=2),
encoding="utf-8")
def next_chunk(today):
"""下一个待回灌窗口 (start, end), newest-first 逐片向 2020 走; None=回灌完成。
片宽 CHUNK_DAYS, 尾片与 FULL_START 对齐(前段有少量重叠, OR REPLACE 去重)。
"""
state = load_chunk_state()
if state.get("next_end") is None:
return None
end = min(dt.date.fromisoformat(state["next_end"]), today)
if end <= dt.date.fromisoformat(FULL_START):
return None
start = max(dt.date.fromisoformat(FULL_START),
end - dt.timedelta(days=CHUNK_DAYS))
return start.isoformat(), end.isoformat()
def _has_unique_index(conn):
"""dbbardata 是否已有覆盖 (symbol,exchange,interval,datetime) 的唯一索引。
NAS 副本主库已带 uq_dbbardata(merge_increment 建)——同列重复建唯一索引 =
数 GB 无谓开销+双倍写放大, 必须识别已有索引跳过(索引名不同语义同)。
"""
for idx in conn.execute("PRAGMA index_list('dbbardata')").fetchall():
if not idx[2]: # (seq, name, unique, origin, partial) 第 3 列 unique
continue
cols = {r[2] for r in conn.execute(
"PRAGMA index_info('%s')" % idx[1]).fetchall()}
if cols == {"symbol", "exchange", "interval", "datetime"}:
return True
return False
def ensure_schema(conn):
"""同库写入前的 schema 保障: 建表(若空库)+唯一索引(若尚无, REPLACE 去重前提)。"""
conn.execute(
"CREATE TABLE IF NOT EXISTS dbbardata ("
"symbol TEXT NOT NULL, exchange TEXT NOT NULL, datetime TEXT NOT NULL, "
"interval TEXT NOT NULL, volume REAL NOT NULL, turnover REAL NOT NULL, "
"open_interest REAL NOT NULL, open_price REAL NOT NULL, "
"high_price REAL NOT NULL, low_price REAL NOT NULL, close_price REAL NOT NULL)"
)
if not _has_unique_index(conn):
conn.execute(
"CREATE UNIQUE INDEX "
"dbbardata_symbol_exchange_interval_datetime "
"ON dbbardata (symbol, exchange, interval, datetime)"
)
conn.commit()
def upsert_5m(conn, code, prefix, rows):
"""5min rows -> dbbardata('5m')。datetime 拼接/数值口径与 bs_eod 15m 完全同款。"""
if not rows:
return 0
df = pd.DataFrame(rows, columns=M5_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": _build_15m_dt(df["date"], df["time"]),
"interval": "5m", "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(conn, code, prefix, start, end, timeout):
"""单股一个短事务(fetch+upsert), 异常整股 rollback —— per-stock 原子断点续跑。"""
bs_code = f"{prefix}.{code}"
with conn:
rows = fetch_k_with_timeout(bs_code, M5_FIELDS, "5", start, end,
timeout=timeout)
return upsert_5m(conn, code, prefix, rows)
def _lock_path():
return DB.parent / "bs_5m.lock"
def _acquire_lock():
"""单实例锁(pid 活性检测): 活 pid 持锁→False; 死 pid/无锁→取锁 True。"""
lock = _lock_path()
if lock.exists():
try:
pid = int(lock.read_text(encoding="utf-8").strip())
os.kill(pid, 0) # 活进程静默返回; 死进程 ProcessLookupError
log.warning("实例 pid=%d 在跑, 本次让路", pid)
return False
except ValueError:
log.warning("锁内容非 pid, 视为陈旧锁覆盖")
except ProcessLookupError:
log.warning("锁持进程已死, 覆盖陈旧锁")
except PermissionError:
return False # 进程存在但属他人 → 按"在跑"处理
DB.parent.mkdir(parents=True, exist_ok=True)
lock.write_text(str(os.getpid()), encoding="utf-8")
return True
def _release_lock():
try:
_lock_path().unlink(missing_ok=True)
except OSError as e:
log.warning("释放锁失败: %s", e)
def main():
ap = argparse.ArgumentParser()
ap.add_argument("--limit", type=int, default=0, help="限股数(冒烟用)")
args = ap.parse_args()
if not _acquire_lock():
log.info("[SKIP] 另一实例在跑, exit 0")
sys.exit(0)
try:
_run(args)
finally:
_release_lock()
def _sweep(conn, stocks, start, end, swept, label):
"""扫全 A 一个窗口: per-stock 短事务 + 周期 relogin + 间隔。
swept 为跨窗口累计股数(=query 数, 每股恰 1 次调用), 由调用方持有并回填;
返 (stats, limit_reached)。
"""
stats = {"ok": 0, "empty": 0, "failed": 0, "db_rows": 0}
limit_reached = False
t0 = time.time()
for i, (code, prefix) in enumerate(stocks):
if swept[0] >= DAILY_LIMIT:
log.warning("[%s] query %d 达防线 %d, graceful stop",
label, swept[0], DAILY_LIMIT)
limit_reached = True
break
swept[0] += 1
try:
n = _process_one(conn, code, prefix, start, end, 60)
stats["db_rows"] += n
if n:
stats["ok"] += 1
else:
stats["empty"] += 1
except Exception as e:
stats["failed"] += 1
if stats["failed"] <= 5 or stats["failed"] % 100 == 0:
log.warning("[%s] %s err: %s", label, code, e)
if not relogin():
log.error("[%s] %s relogin 失败, 跳过", label, code)
if (i + 1) % RELOGIN_EVERY == 0:
log.info("[%s] 进度 %d/%d ok=%d empty=%d failed=%d rows=%d (%.0fs)",
label, i + 1, len(stocks), stats["ok"], stats["empty"],
stats["failed"], stats["db_rows"], time.time() - t0)
if not relogin():
log.warning("[%s] 周期 relogin 失败, 继续跑", label)
if i < len(stocks) - 1:
time.sleep(BS_INTERVAL)
return stats, limit_reached
def _run(args):
today = dt.date.today()
chunk = next_chunk(today)
inc_start = (today - dt.timedelta(days=LOOKBACK)).strftime("%Y-%m-%d")
inc_end = today.strftime("%Y-%m-%d")
log.info("bs_5m start db=%s chunk=%s inc=%s~%s", DB,
"%s~%s" % chunk if chunk else "无(回灌完成)", inc_start, inc_end)
if not login_with_retry():
log.error("[SKIP] 登录失败, 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]
conn = sqlite3.connect(str(DB), timeout=60)
conn.execute("PRAGMA busy_timeout = 60000")
conn.execute("PRAGMA journal_mode = WAL")
ensure_schema(conn)
swept = [0] # 跨窗口累计 query 数(每股恰 1 次调用)
limit_reached = False
t0 = time.time()
try:
if chunk:
cs, lr = _sweep(conn, stocks, chunk[0], chunk[1], swept,
"chunk %s~%s" % chunk)
limit_reached = limit_reached or lr
log.info("[CHUNK] %s~%s ok=%d empty=%d failed=%d rows=%d",
chunk[0], chunk[1], cs["ok"], cs["empty"], cs["failed"],
cs["db_rows"])
if not lr and swept[0] < DAILY_LIMIT:
# 整片扫完才推进 state; 没扫完次日重拉同片(OR REPLACE 幂等续跑)
save_chunk_state({"next_end": chunk[0]})
# 每日增量照跑(回灌期也跑, 近端 7 天始终新鲜)
istats, lr = _sweep(conn, stocks, inc_start, inc_end, swept, "inc")
limit_reached = limit_reached or lr
log.info("[INC] ok=%d empty=%d failed=%d rows=%d",
istats["ok"], istats["empty"], istats["failed"],
istats["db_rows"])
if next_chunk(today) is None and not limit_reached:
save_chunk_state({"next_end": None}) # 回灌完结盖章(幂等)
finally:
conn.close()
try:
import baostock as bs
bs.logout()
except Exception:
pass
log.info("[DONE] query=%d 耗时%.0fs", swept[0], time.time() - t0)
sys.exit(3 if limit_reached else 0)
if __name__ == "__main__":
main()