diff --git a/scripts/data_platform/bs_5m_eod.py b/scripts/data_platform/bs_5m_eod.py index 3246957..76cac89 100644 --- a/scripts/data_platform/bs_5m_eod.py +++ b/scripts/data_platform/bs_5m_eod.py @@ -1,14 +1,22 @@ #!/usr/bin/env python3 # -*- coding: utf-8 -*- -"""bs_5m_eod.py — NAS 独立 5 分钟线备份 (数据补全 2026-08-20 用户拍板). +"""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 等"以后有机会"。 +定位: 5m 直接写进 NAS 副本主库 dbbardata 表(与 VPS 同步镜像同一个 db/同一张表), +靠 interval='5m' 与镜像行(d/15m)在唯一键上正交实现互不干扰 —— 用户目标: NAS 就地 +5m 回测/回放(reader 零改动), 未来 VPS 扩容直接把 interval='5m' 导出 merge 过去。 -- 库: BS_5M_DB 缺省 /volume1/stock/sanguo_5m/dbbardata_5m.db(容器内外同路径挂载, - 首次自建 schema+唯一索引) +同库安全性(结构性, 非约定): +- 同步链(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 短事务, 中断重跑幂等) @@ -35,7 +43,8 @@ from bs_eod import ( # noqa: E402 — 复用同一套健壮性封装, 口径与 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")) +_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 超时不够 @@ -56,8 +65,24 @@ def calc_window(today, full): 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(与生产 dbbardata 同列) + 唯一索引(REPLACE 去重前提)。""" + """同库写入前的 schema 保障: 建表(若空库)+唯一索引(若尚无, REPLACE 去重前提)。""" conn.execute( "CREATE TABLE IF NOT EXISTS dbbardata (" "symbol TEXT NOT NULL, exchange TEXT NOT NULL, datetime TEXT NOT NULL, " @@ -65,11 +90,12 @@ def ensure_schema(conn): "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)" - ) + if not _has_unique_index(conn): + conn.execute( + "CREATE UNIQUE INDEX " + "dbbardata_symbol_exchange_interval_datetime " + "ON dbbardata (symbol, exchange, interval, datetime)" + ) conn.commit() diff --git a/scripts/nas_sync/merge_increment.py b/scripts/nas_sync/merge_increment.py index 92bd30f..18a90ac 100644 --- a/scripts/nas_sync/merge_increment.py +++ b/scripts/nas_sync/merge_increment.py @@ -31,7 +31,10 @@ def main(): conn = sqlite3.connect(args.db) cur = conn.cursor() # 大表(1.8亿行+)merge 提速: 默认2MB cache 致 UNIQUE索引(数GB)全磁盘IO, merge>120s超时 + # busy_timeout(2026-08-20): NAS 副本从此有第二个写者(bs_5m_eod 5分钟线), + # 无 busy_timeout 撞写锁会当场 "database is locked" 失败; 60s 与 5m 侧对齐。 for _p in ("PRAGMA journal_mode=WAL", "PRAGMA synchronous=NORMAL", + "PRAGMA busy_timeout=60000", "PRAGMA cache_size=-500000", "PRAGMA temp_store=MEMORY"): cur.execute(_p) cur.execute(SCHEMA) diff --git a/tests/data_platform/test_bs_5m_eod.py b/tests/data_platform/test_bs_5m_eod.py index ecc54cc..7e100d1 100644 --- a/tests/data_platform/test_bs_5m_eod.py +++ b/tests/data_platform/test_bs_5m_eod.py @@ -1,9 +1,10 @@ # -*- coding: utf-8 -*- -"""TDD for bs_5m_eod.py — NAS 独立 5 分钟线备份 (2026-08-20 用户拍板). +"""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。 +定位: 5m 直接写 NAS 副本主库 dbbardata 表(与 VPS 同步镜像同库同表), 靠 +interval='5m' 与镜像行在唯一键上正交互不干扰 —— 同步链导出不带 id + NAS 侧 +INSERT OR IGNORE 去重(merge_increment.py), NAS 回测 reader 零改动可见 5m, +未来 VPS 扩容导 interval='5m' 反向 merge 即迁移。 配额: NAS 出口 IP 独立核算 48k/天, 一趟仅 ~5.5k query(每股恰好 1 次调用, --full 全区间也一样); VPS 侧 bs_eod/bs_fund 用 VPS 自己的 IP 配额, 互不相干。 @@ -51,6 +52,12 @@ def _m5_rows(bs_code, n=1, hhmm="0935"): # ---------- schema ---------- +def test_default_db_is_nas_replica_main_db(): + """契约钉死: 缺省库 = NAS 副本主库(同步目标=回测源), 绝非旁路独立库。""" + assert b5._DEFAULT_DB == \ + "/volume1/stock/sanguo_vnpy_v2/data_backup/quant_trading.db" + + def test_ensure_schema_creates_table_and_unique_index(tmp_path): db = tmp_path / "x.db" conn = sqlite3.connect(str(db)) @@ -62,6 +69,29 @@ def test_ensure_schema_creates_table_and_unique_index(tmp_path): conn.close() +def test_ensure_schema_skips_duplicate_index_on_replica(tmp_path): + """NAS 副本主库已带 uq_dbbardata(merge_increment 建, 名字不同语义同): + ensure_schema 不得再建同列唯一索引(数 GB 无谓开销+双倍写放大)。""" + db = tmp_path / "replica.db" + conn = sqlite3.connect(str(db)) + conn.execute( + "CREATE TABLE dbbardata (id INTEGER PRIMARY KEY AUTOINCREMENT, " + "symbol TEXT, exchange TEXT, datetime TEXT, interval TEXT, " + "volume REAL, turnover REAL, open_interest REAL, " + "open_price REAL, high_price REAL, low_price REAL, close_price REAL)") + conn.execute( + "CREATE UNIQUE INDEX uq_dbbardata " + "ON dbbardata (symbol, exchange, interval, datetime)") + conn.commit() + b5.ensure_schema(conn) + unique_idx = [r for r in conn.execute( + "SELECT name FROM sqlite_master " + "WHERE tbl_name='dbbardata' AND type='index' AND sql IS NOT NULL")] + assert len(unique_idx) == 1 # 只有原有的 uq_dbbardata, 没有重复建 + assert unique_idx[0][0] == "uq_dbbardata" + conn.close() + + # ---------- upsert_5m ---------- def test_upsert_5m_writes_dbbardata(tmp_db_conn):