311 lines
12 KiB
Python
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()
|