feat(data): NAS独立5分钟线备份——bs_5m_eod新脚本(NAS专属,VPS永不跑)——用户拍板2026-08-20:NAS既有数据全是VPS同步镜像(nas_sync零接触),5m是NAS唯一自主从baostock下载数据,落独立库BS_5M_DB(缺省/volume1/stock/sanguo_5m/dbbardata_5m.db容器内外同路径),绝不写同步目标data_backup/quant_trading.db;①首次自建schema+唯一索引(REPLACE去重)②缺省每日增量LOOKBACK=7(DSM任务计划~19:05);--full一次性回灌2020-01-03+(baostock分钟固定起点非滚动;全区间每股恰1次调用=~5.5k query一趟≈3.6亿行,数小时级,per-stock短事务断点续跑幂等;FULL_TIMEOUT=300s防大payload误杀)③配额:NAS出口IP独立核算48k/天与VPS互不相干,DAILY_LIMIT守卫保留④单实例锁bs_5m.lock(pid活性检测,死锁自动覆盖/PermissionError按在跑)防DSM定时与--full长跑撞车=同IP双baostock连接红线⑤复用bs_eod四件套+datetime拼接/数值口径与15m完全同款;+13测试(schema/REPLACE不双行/per-stock原子/full vs daily窗口/锁三态/main四条);data_platform 153绿 [nas]
This commit is contained in:
@@ -0,0 +1,225 @@
|
||||
#!/usr/bin/env python3
|
||||
# -*- coding: utf-8 -*-
|
||||
"""bs_5m_eod.py — NAS 独立 5 分钟线备份 (数据补全 2026-08-20 用户拍板).
|
||||
|
||||
定位(铁律): NAS 的既有数据全部是 VPS 同步镜像(nas_sync 机制, 本脚本零接触);
|
||||
5m 是 NAS 唯一自主从 baostock 下载的数据, 落独立库 BS_5M_DB —— 绝不写
|
||||
/volume1/stock/sanguo_vnpy_v2/data_backup/quant_trading.db(VPS 同步目标, 写了
|
||||
下次同步即被镜像覆盖/冲突)。VPS 不跑本脚本(无 schtask), 5m 上 VPS 等"以后有机会"。
|
||||
|
||||
- 库: BS_5M_DB 缺省 /volume1/stock/sanguo_5m/dbbardata_5m.db(容器内外同路径挂载,
|
||||
首次自建 schema+唯一索引)
|
||||
- 模式: 缺省每日增量(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,
|
||||
)
|
||||
|
||||
DB = Path(os.environ.get("BS_5M_DB", "/volume1/stock/sanguo_5m/dbbardata_5m.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 ensure_schema(conn):
|
||||
"""独立库自建 schema(与生产 dbbardata 同列) + 唯一索引(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)"
|
||||
)
|
||||
conn.execute(
|
||||
"CREATE UNIQUE INDEX IF NOT EXISTS "
|
||||
"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()
|
||||
@@ -0,0 +1,189 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""TDD for bs_5m_eod.py — NAS 独立 5 分钟线备份 (2026-08-20 用户拍板).
|
||||
|
||||
定位: NAS 自己从 baostock 下载 5m, 不推 VPS; NAS 其余数据全是 VPS 同步镜像,
|
||||
本脚本是 NAS 唯一自主数据 —— 落独立库 BS_5M_DB(缺省 /volume1/stock/sanguo_5m/),
|
||||
绝不碰同步目录 /volume1/stock/sanguo_vnpy_v2/data_backup/quant_trading.db。
|
||||
|
||||
配额: NAS 出口 IP 独立核算 48k/天, 一趟仅 ~5.5k query(每股恰好 1 次调用,
|
||||
--full 全区间也一样); VPS 侧 bs_eod/bs_fund 用 VPS 自己的 IP 配额, 互不相干。
|
||||
"""
|
||||
import datetime as dt
|
||||
import os
|
||||
import sqlite3
|
||||
import sys
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
if "baostock" not in sys.modules:
|
||||
sys.modules["baostock"] = MagicMock()
|
||||
|
||||
from scripts.data_platform import bs_5m_eod as b5 # noqa: E402
|
||||
|
||||
|
||||
# ---------- Fixtures ----------
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _fast(monkeypatch):
|
||||
monkeypatch.setattr(b5, "BS_INTERVAL", 0.0)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def tmp_db_conn(tmp_path, monkeypatch):
|
||||
"""重定向 DB 到 tmp_path, 返已建 schema 的连接(测完关)."""
|
||||
db = tmp_path / "dbbardata_5m.db"
|
||||
monkeypatch.setattr(b5, "DB", db)
|
||||
conn = sqlite3.connect(str(db))
|
||||
b5.ensure_schema(conn)
|
||||
yield conn
|
||||
conn.close()
|
||||
|
||||
|
||||
def _m5_rows(bs_code, n=1, hhmm="0935"):
|
||||
"""模拟 baostock 5min row (9 列 M5_FIELDS, time 17 位)."""
|
||||
out = []
|
||||
for i in range(n):
|
||||
out.append(["2026-07-01", f"20260701{hhmm}00000", bs_code,
|
||||
"10", "11", "9", "10.5", "1000", "10000"])
|
||||
return out
|
||||
|
||||
|
||||
# ---------- schema ----------
|
||||
|
||||
def test_ensure_schema_creates_table_and_unique_index(tmp_path):
|
||||
db = tmp_path / "x.db"
|
||||
conn = sqlite3.connect(str(db))
|
||||
b5.ensure_schema(conn)
|
||||
names = {r[0] for r in conn.execute(
|
||||
"SELECT name FROM sqlite_master WHERE tbl_name='dbbardata'")}
|
||||
assert "dbbardata" in names
|
||||
assert "dbbardata_symbol_exchange_interval_datetime" in names # REPLACE 去重前提
|
||||
conn.close()
|
||||
|
||||
|
||||
# ---------- upsert_5m ----------
|
||||
|
||||
def test_upsert_5m_writes_dbbardata(tmp_db_conn):
|
||||
n = b5.upsert_5m(tmp_db_conn, "000001", "sz", _m5_rows("sz.000001"))
|
||||
tmp_db_conn.commit()
|
||||
assert n == 1
|
||||
row = tmp_db_conn.execute(
|
||||
"SELECT symbol, exchange, datetime, interval, close_price FROM dbbardata"
|
||||
).fetchone()
|
||||
assert row[0] == "000001"
|
||||
assert row[1] == "SZSE"
|
||||
assert row[2] == "2026-07-01 09:35:00" # 17 位 time 第 8-12 位取 HHMM
|
||||
assert row[3] == "5m"
|
||||
assert abs(row[4] - 10.5) < 1e-6
|
||||
|
||||
|
||||
def test_upsert_5m_replaces_same_pk_no_dup(tmp_db_conn):
|
||||
b5.upsert_5m(tmp_db_conn, "000001", "sz", _m5_rows("sz.000001"))
|
||||
b5.upsert_5m(tmp_db_conn, "000001", "sz", _m5_rows("sz.000001"))
|
||||
tmp_db_conn.commit()
|
||||
assert tmp_db_conn.execute("SELECT COUNT(*) FROM dbbardata").fetchone()[0] == 1
|
||||
|
||||
|
||||
def test_upsert_5m_empty_rows_no_op(tmp_db_conn):
|
||||
assert b5.upsert_5m(tmp_db_conn, "000001", "sz", []) == 0
|
||||
assert tmp_db_conn.execute("SELECT COUNT(*) FROM dbbardata").fetchone()[0] == 0
|
||||
|
||||
|
||||
def test_process_one_per_stock_atomic(tmp_db_conn, monkeypatch):
|
||||
"""单股独立事务: fetch 抛错 → 整股 rollback 不留半行(per-stock 原子, house style)."""
|
||||
with patch.object(b5, "fetch_k_with_timeout",
|
||||
side_effect=RuntimeError("baostock hiccup")):
|
||||
with pytest.raises(RuntimeError):
|
||||
b5._process_one(tmp_db_conn, "000001", "sz", "2026-08-01",
|
||||
"2026-08-19", timeout=60)
|
||||
assert tmp_db_conn.execute("SELECT COUNT(*) FROM dbbardata").fetchone()[0] == 0
|
||||
|
||||
|
||||
def test_process_one_success_writes(tmp_db_conn):
|
||||
with patch.object(b5, "fetch_k_with_timeout",
|
||||
side_effect=lambda *a, **k: _m5_rows("sz.000001")):
|
||||
n = b5._process_one(tmp_db_conn, "000001", "sz", "2026-08-01",
|
||||
"2026-08-19", timeout=60)
|
||||
assert n == 1
|
||||
|
||||
|
||||
# ---------- 窗口 ----------
|
||||
|
||||
def test_calc_window_full_vs_daily():
|
||||
today = dt.date(2026, 8, 20)
|
||||
assert b5.calc_window(today, full=True) == ("2020-01-03", "2026-08-20")
|
||||
s, e = b5.calc_window(today, full=False)
|
||||
assert s == "2026-08-13" and e == "2026-08-20" # LOOKBACK=7
|
||||
|
||||
|
||||
# ---------- 单实例锁(防 DSM 定时与全量回灌撞车) ----------
|
||||
|
||||
def test_lock_acquire_and_second_refused(tmp_path, monkeypatch):
|
||||
monkeypatch.setattr(b5, "DB", tmp_path / "dbbardata_5m.db")
|
||||
assert b5._acquire_lock() is True # 活进程(自己)持锁
|
||||
assert b5._acquire_lock() is False # 第二实例拒绝
|
||||
b5._release_lock()
|
||||
assert b5._acquire_lock() is True # 释放后可再取
|
||||
|
||||
|
||||
def test_lock_stale_pid_reclaimed(tmp_path, monkeypatch):
|
||||
monkeypatch.setattr(b5, "DB", tmp_path / "dbbardata_5m.db")
|
||||
lock = tmp_path / "bs_5m.lock"
|
||||
lock.write_text("999999", encoding="utf-8") # 死 pid(无此进程)
|
||||
assert b5._acquire_lock() is True
|
||||
b5._release_lock()
|
||||
|
||||
|
||||
# ---------- main 集成 ----------
|
||||
|
||||
@pytest.fixture
|
||||
def main_env(tmp_path, monkeypatch):
|
||||
monkeypatch.setattr(b5, "DB", tmp_path / "dbbardata_5m.db")
|
||||
monkeypatch.setattr(sys, "argv", ["bs_5m_eod.py"])
|
||||
monkeypatch.setattr(b5, "login_with_retry", MagicMock(return_value=True))
|
||||
monkeypatch.setattr(b5, "fetch_all_stocks_with_timeout",
|
||||
MagicMock(return_value=[("600001", "sh"), ("000002", "sz")]))
|
||||
return tmp_path
|
||||
|
||||
|
||||
def test_main_runs_and_exits_0(main_env):
|
||||
with patch.object(b5, "_process_one", return_value=48) as m_proc:
|
||||
with pytest.raises(SystemExit) as e:
|
||||
b5.main()
|
||||
assert e.value.code == 0
|
||||
assert m_proc.call_count == 2
|
||||
# 数据库文件+schema 已就位
|
||||
conn = sqlite3.connect(str(b5.DB))
|
||||
assert conn.execute(
|
||||
"SELECT COUNT(*) FROM sqlite_master "
|
||||
"WHERE tbl_name='dbbardata' AND type='table'"
|
||||
).fetchone()[0] == 1
|
||||
conn.close()
|
||||
# 正常退出释放锁
|
||||
assert not (main_env / "bs_5m.lock").exists()
|
||||
|
||||
|
||||
def test_main_exit2_when_login_fails(main_env, monkeypatch):
|
||||
monkeypatch.setattr(b5, "login_with_retry", MagicMock(return_value=False))
|
||||
with pytest.raises(SystemExit) as e:
|
||||
b5.main()
|
||||
assert e.value.code == 2
|
||||
|
||||
|
||||
def test_main_exit3_when_limit_reached(main_env, monkeypatch):
|
||||
monkeypatch.setattr(b5, "DAILY_LIMIT", 0) # 每股恰 1 query, 0=立即达限
|
||||
with patch.object(b5, "_process_one",
|
||||
side_effect=AssertionError("不应跑股")):
|
||||
with pytest.raises(SystemExit) as e:
|
||||
b5.main()
|
||||
assert e.value.code == 3
|
||||
|
||||
|
||||
def test_main_second_instance_skips(main_env, monkeypatch):
|
||||
monkeypatch.setattr(b5, "_acquire_lock", MagicMock(return_value=False))
|
||||
monkeypatch.setattr(b5, "login_with_retry",
|
||||
MagicMock(side_effect=AssertionError("不应登录")))
|
||||
with pytest.raises(SystemExit) as e:
|
||||
b5.main()
|
||||
assert e.value.code == 0
|
||||
Reference in New Issue
Block a user