Files
sanguo_vnpy_v2/scripts/data_platform/bs_5m_eod.py
T

252 lines
10 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)
则不重复建
- 模式: 缺省每日增量(LOOKBACK 7 天, DSM 任务计划 ~19:05); --full 一次性回灌
2020-01-03+(baostock 分钟数据固定起点非滚动; 全区间每股恰好 1 次调用, 一趟
~5.5k query / ~3.6 亿行, 数小时级; per-stock 短事务, 中断重跑幂等)
- 配额: NAS 出口 IP 独立核算 48k/天(与 VPS 各自计数互不相干); 本脚本一趟远低
于限, DAILY_LIMIT 守卫保留
- 单实例锁 bs_5m.lock(pid 活性检测): 防 DSM 定时与 --full 长跑撞车
(同 IP 双 baostock 连接红线)
退出码: 0=完成/让路; 1=致命; 2=登录失败; 3=query 超限 graceful stop
"""
import argparse
import datetime as dt
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)
FULL_TIMEOUT = 300 # 全区间单股 ~6.4 万行, 默认 60s 超时不够
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 calc_window(today, full):
"""--full: 2020-01-03 固定起点全区间; 缺省: LOOKBACK 天增量窗口。"""
end = today.strftime("%Y-%m-%d")
if full:
return FULL_START, end
start = (today - dt.timedelta(days=LOOKBACK)).strftime("%Y-%m-%d")
return start, end
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="限股数(冒烟用)")
ap.add_argument("--full", action="store_true",
help="一次性回灌 2020-01-03+(缺省每日增量)")
args = ap.parse_args()
if not _acquire_lock():
log.info("[SKIP] 另一实例在跑, exit 0")
sys.exit(0)
try:
_run(args)
finally:
_release_lock()
def _run(args):
today = dt.date.today()
start, end = calc_window(today, args.full)
log.info("bs_5m start window=%s~%s mode=%s db=%s", start, end,
"full" if args.full else "daily", DB)
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)
timeout = FULL_TIMEOUT if args.full else 60
stats = {"ok": 0, "empty": 0, "failed": 0, "db_rows": 0}
limit_reached = False
t0 = time.time()
try:
for i, (code, prefix) in enumerate(stocks):
if i >= DAILY_LIMIT: # 每股恰好 1 query, i 即当日 query 计数
log.warning("query %d 达防线 %d, graceful stop", i, DAILY_LIMIT)
limit_reached = True
break
try:
n = _process_one(conn, code, prefix, start, end, timeout)
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 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 rows=%d (%.0fs)",
i + 1, len(stocks), stats["ok"], stats["empty"],
stats["failed"], stats["db_rows"], time.time() - t0)
if not relogin():
log.warning("周期 relogin 失败, 继续跑")
if i < len(stocks) - 1:
time.sleep(BS_INTERVAL)
finally:
conn.close()
try:
import baostock as bs
bs.logout()
except Exception:
pass
log.info("[DONE] mode=%s ok=%d empty=%d failed=%d db_rows=%d 耗时%.0fs",
"full" if args.full else "daily", stats["ok"], stats["empty"],
stats["failed"], stats["db_rows"], time.time() - t0)
sys.exit(3 if limit_reached else 0)
if __name__ == "__main__":
main()