Files
sanguo_vnpy_v2/scripts/nas_sync/merge_increment.py
T
claude_dev d76171434b fix(nas_sync): Phase3 full 修复(STEP降片+merge PRAGMA)+全量完成2.02亿行
- sync_dbbardata.sh: STEP 2000万→200万/片(首片2GB export+scp跨公网挂)
- merge_increment.py: 加 PRAGMA cache=500MB+WAL+NORMAL+synchronous
  (replica 2亿行 UNIQUE索引数GB,默认2MB cache致索引IO merge>120s超时)
- 全量完成: replica 2.019亿行(=dbbardata全量日线+15min+ETF)
- 踩坑链: ssh reset中断→timeout杀ssh留NAS merge残留持锁→since落后数据已在库
- 教训: timeout杀ssh不杀远程python;merge大表必加PRAGMA;断点续传since.txt救命

Co-Authored-By: Claude <noreply@anthropic.com>
2026-07-30 18:08:42 +08:00

60 lines
2.4 KiB
Python

"""NAS 端:把增量 sqlite 合并进本地 dbbardata 副本。
设计要点:
- 纯 sqlite3 标准库,零第三方依赖(NAS 宿主 Python 3.8 也能跑)。
- INSERT OR IGNORE ... SELECT:靠 UNIQUE(symbol,exchange,interval,datetime)
去重,重复行安全跳过(增量重跑、首全量分片重叠都不怕)。
- 首次合并自动建表+建唯一索引(与 VPS schema 对齐,id AUTOINCREMENT)。
用法:
python merge_increment.py --db /volume1/.../quant_trading.db --inc /tmp/inc.db
"""
import argparse
import sqlite3
SCHEMA = """CREATE TABLE IF NOT EXISTS 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)"""
COLS = ("symbol,exchange,datetime,interval,volume,turnover,open_interest,"
"open_price,high_price,low_price,close_price")
def main():
ap = argparse.ArgumentParser(description="合并增量 sqlite 进本地 dbbardata 副本")
ap.add_argument("--db", required=True, help="NAS 本地 quant_trading.db 路径")
ap.add_argument("--inc", required=True, help="增量 sqlite 路径")
args = ap.parse_args()
conn = sqlite3.connect(args.db)
cur = conn.cursor()
# 大表(1.8亿行+)merge 提速: 默认2MB cache 致 UNIQUE索引(数GB)全磁盘IO, merge>120s超时
for _p in ("PRAGMA journal_mode=WAL", "PRAGMA synchronous=NORMAL",
"PRAGMA cache_size=-500000", "PRAGMA temp_store=MEMORY"):
cur.execute(_p)
cur.execute(SCHEMA)
cur.execute(
"CREATE UNIQUE INDEX IF NOT EXISTS uq_dbbardata "
"ON dbbardata(symbol,exchange,interval,datetime)")
cur.execute("ATTACH DATABASE ? AS inc", (args.inc,))
inc_count = cur.execute("SELECT COUNT(*) FROM inc.dbbardata").fetchone()[0]
before = cur.execute("SELECT COUNT(*) FROM main.dbbardata").fetchone()[0]
cur.execute(
"INSERT OR IGNORE INTO main.dbbardata(%s) SELECT %s FROM inc.dbbardata"
% (COLS, COLS))
conn.commit()
after = cur.execute("SELECT COUNT(*) FROM main.dbbardata").fetchone()[0]
conn.close()
inserted = after - before
skipped = inc_count - inserted
print("MERGED inc=%d inserted=%d skipped(dup)=%d before=%d after=%d"
% (inc_count, inserted, skipped, before, after))
if __name__ == "__main__":
main()