From 20ff70e85db6bac82b9a364be95ad76217754138 Mon Sep 17 00:00:00 2001 From: claude_dev Date: Thu, 20 Aug 2026 10:10:04 +0800 Subject: [PATCH] =?UTF-8?q?refactor(data):=205m=E6=94=B9=E5=90=8C=E5=BA=93?= =?UTF-8?q?=E8=AE=BE=E8=AE=A1=E2=80=94=E2=80=94=E7=94=A8=E6=88=B7=E4=BA=8C?= =?UTF-8?q?=E6=AC=A1=E6=8B=8D=E6=9D=BF2026-08-20:=E6=94=BE=E5=BC=83?= =?UTF-8?q?=E7=8B=AC=E7=AB=8B=E5=BA=93,5m=E7=9B=B4=E6=8E=A5=E5=86=99NAS?= =?UTF-8?q?=E5=89=AF=E6=9C=AC=E4=B8=BB=E5=BA=93dbbardata=E8=A1=A8(interval?= =?UTF-8?q?=3D'5m'=E4=B8=8E=E9=95=9C=E5=83=8F=E8=A1=8Cd/15m=E5=94=AF?= =?UTF-8?q?=E4=B8=80=E9=94=AE=E6=AD=A3=E4=BA=A4=3D=E7=BB=93=E6=9E=84?= =?UTF-8?q?=E6=80=A7=E4=BA=92=E4=B8=8D=E5=B9=B2=E6=89=B0)=E2=80=94?= =?UTF-8?q?=E2=80=94=E2=91=A0bs=5F5m=5Feod=E7=BC=BA=E7=9C=81=E5=BA=93?= =?UTF-8?q?=E6=94=B9/volume1/stock/sanguo=5Fvnpy=5Fv2/data=5Fbackup/quant?= =?UTF-8?q?=5Ftrading.db(=3D=E5=90=8C=E6=AD=A5=E7=9B=AE=E6=A0=87=3DNAS?= =?UTF-8?q?=E5=9B=9E=E6=B5=8Bprovider=E8=AF=BB=E7=9A=84=E5=BA=93portfolio?= =?UTF-8?q?=5Fworker.py:180),NAS=E5=B0=B1=E5=9C=B05m=E5=9B=9E=E6=B5=8B/?= =?UTF-8?q?=E5=9B=9E=E6=94=BEreader=E9=9B=B6=E6=94=B9=E5=8A=A8,=E6=9C=AA?= =?UTF-8?q?=E6=9D=A5VPS=E6=89=A9=E5=AE=B9=E5=AF=BCinterval=3D'5m'=E5=8F=8D?= =?UTF-8?q?=E5=90=91merge=E4=B8=80=E6=AC=A1=E8=BF=81=E7=A7=BB=E2=91=A1?= =?UTF-8?q?=E5=90=8C=E5=BA=93=E5=AE=89=E5=85=A8=E6=80=A7=E5=AE=9E=E8=AF=81?= =?UTF-8?q?:=E5=90=8C=E6=AD=A5=E9=93=BEexport=E4=B8=8D=E5=B8=A6id=E5=88=97?= =?UTF-8?q?+NAS=E4=BE=A7INSERT=20OR=20IGNORE=E6=8C=89UNIQUE=E5=8E=BB?= =?UTF-8?q?=E9=87=8D=E2=86=92=E9=94=AE=E6=B0=B8=E4=B8=8D=E7=9B=B8=E4=BA=A4?= =?UTF-8?q?=E5=8F=8C=E5=90=91=E7=A2=B0=E4=B8=8D=E5=88=B0;=E7=BA=AFmerge?= =?UTF-8?q?=E6=97=A0=E6=96=87=E4=BB=B6=E8=A6=86=E7=9B=96NAS=E6=9C=AC?= =?UTF-8?q?=E5=9C=B0=E8=A1=8C=E4=B8=8D=E4=BC=9A=E8=A2=AB=E5=88=A0;id?= =?UTF-8?q?=E7=94=B1NAS=20AUTOINCREMENT=E8=87=AA=E5=88=86=E9=85=8D?= =?UTF-8?q?=E6=97=A0=E8=B7=A8=E5=BA=93=E5=86=B2=E7=AA=81=E2=91=A2ensure=5F?= =?UTF-8?q?schema=E9=98=B2=E9=87=8D=E5=A4=8D=E7=B4=A2=E5=BC=95:PRAGMA=20in?= =?UTF-8?q?dex=5Flist/index=5Finfo=E8=AF=86=E5=88=AB=E5=B7=B2=E6=9C=89?= =?UTF-8?q?=E5=90=8C=E5=88=97=E5=94=AF=E4=B8=80=E7=B4=A2=E5=BC=95(NAS?= =?UTF-8?q?=E4=BE=A7uq=5Fdbbardata)=E8=B7=B3=E8=BF=87,=E5=85=8D=E6=95=B0GB?= =?UTF-8?q?=E6=97=A0=E8=B0=93=E5=BC=80=E9=94=80+=E5=8F=8C=E5=80=8D?= =?UTF-8?q?=E5=86=99=E6=94=BE=E5=A4=A7=E2=91=A3merge=5Fincrement=E8=A1=A5b?= =?UTF-8?q?usy=5Ftimeout=3D60000(=E5=89=AF=E6=9C=AC=E5=BA=93=E4=BB=8E?= =?UTF-8?q?=E6=AD=A4=E6=9C=89=E7=AC=AC=E4=BA=8C=E4=B8=AA=E5=86=99=E8=80=85?= =?UTF-8?q?,=E6=97=A0=E5=AE=83=E6=92=9E=E5=86=99=E9=94=81=E5=BD=93?= =?UTF-8?q?=E5=9C=BAlocked=E5=A4=B1=E8=B4=A5,60s=E4=B8=8E5m=E4=BE=A7?= =?UTF-8?q?=E5=AF=B9=E9=BD=90)=E2=91=A4=E5=86=99=E9=94=81=E7=AB=9E?= =?UTF-8?q?=E4=BA=89=E8=AF=B4=E6=98=8E:=E5=8F=8C=E5=86=99=E8=80=85WAL+busy?= =?UTF-8?q?=5Ftimeout=E4=B8=B2=E8=A1=8C=E5=8C=96,=E4=B8=8E=E6=AF=8F?= =?UTF-8?q?=E6=97=A5=E5=A2=9E=E9=87=8Fmerge(=E7=A7=92=E7=BA=A7)=E9=87=8D?= =?UTF-8?q?=E5=8F=A0=E5=8D=95=E8=82=A1=E5=A4=B1=E8=B4=A5=E6=AC=A1=E6=97=A5?= =?UTF-8?q?LOOKBACK=3D7=E8=87=AA=E6=84=88;+2=E6=B5=8B=E8=AF=95(=E7=BC=BA?= =?UTF-8?q?=E7=9C=81=E5=BA=93=E5=A5=91=E7=BA=A6=E9=92=89=E6=AD=BB/?= =?UTF-8?q?=E5=89=AF=E6=9C=AC=E5=B7=B2=E6=9C=89=E7=B4=A2=E5=BC=95=E4=B8=8D?= =?UTF-8?q?=E9=87=8D=E5=A4=8D=E5=BB=BA);40+153=E7=BB=BF(data=5Fplatform=20?= =?UTF-8?q?155)=20[nas]?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/data_platform/bs_5m_eod.py | 54 ++++++++++++++++++++------- scripts/nas_sync/merge_increment.py | 3 ++ tests/data_platform/test_bs_5m_eod.py | 38 +++++++++++++++++-- 3 files changed, 77 insertions(+), 18 deletions(-) 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):