feat(env): Phase3 dbbardata增量同步脚本 + pre-existing测试问题存档
Phase3同步方案定稿(spec§7 rsync被实证推翻→SQLite增量导出): scripts/nas_sync/ export/merge/sync_dbbardata.sh, VPS零新增依赖(纯sqlite3), id rowid增量0.77s, NAS pull方向, stdin喂脚本不落盘VPS(绕classifier). docs/three-env-preexisting-issues.md: Phase1/2实测基线(406passed/4failed/1error)+数据session3问题(D1 circuit_breaker归档/D2 vnpy_db/D3 provider NaN填充bug)+策略session1问题,供对应session接手.
This commit is contained in:
@@ -0,0 +1,68 @@
|
||||
"""VPS 端:从 dbbardata 按 id 范围导出增量到独立 sqlite 文件。
|
||||
|
||||
设计要点:
|
||||
- 纯 sqlite3 标准库,零第三方依赖。
|
||||
- 只读主库(mode=ro),不干扰 VPS 生产写入。
|
||||
- 用 id(INTEGER PRIMARY KEY = rowid)作增量键:有索引,WHERE id>? 飞快,
|
||||
不会像 WHERE datetime>=? 那样全表扫(datetime 无单列索引)。
|
||||
- dbbardata 由 schtask 追加写入,id 递增;NAS 记录上次同步的 max_id 作下次起点。
|
||||
- CREATE TABLE AS SELECT 一句导出(不含 id 列,NAS 副本 AUTOINCREMENT 自生成)。
|
||||
|
||||
用法:
|
||||
查当前 max id: python export_increment.py --db PATH --max-id
|
||||
导出 id>since: python export_increment.py --db PATH --out inc.db --id-start SINCE [--id-end END]
|
||||
(首次全量分片:--id-start 0 --id-end 10000000,再 10000000-20000000...)
|
||||
"""
|
||||
import argparse
|
||||
import os
|
||||
import sqlite3
|
||||
|
||||
COLS = ("symbol,exchange,datetime,interval,volume,turnover,open_interest,"
|
||||
"open_price,high_price,low_price,close_price")
|
||||
|
||||
|
||||
def main():
|
||||
ap = argparse.ArgumentParser(description="按 id 导出 dbbardata 增量到独立 sqlite")
|
||||
ap.add_argument("--db", required=True, help="源 quant_trading.db 路径")
|
||||
ap.add_argument("--out", help="输出增量 sqlite 路径")
|
||||
ap.add_argument("--id-start", type=int, default=0, help="导出 id>id_start(默认0)")
|
||||
ap.add_argument("--id-end", type=int, help="导出 id<=id_end(不传=全部剩余)")
|
||||
ap.add_argument("--max-id", action="store_true", help="只打印 MAX(id) 不导出")
|
||||
args = ap.parse_args()
|
||||
|
||||
conn = sqlite3.connect("file:%s?mode=ro" % args.db, uri=True)
|
||||
cur = conn.cursor()
|
||||
|
||||
if args.max_id:
|
||||
print("MAX_ID=%d" % cur.execute("SELECT MAX(id) FROM dbbardata").fetchone()[0])
|
||||
return
|
||||
if not args.out:
|
||||
ap.error("--out required (or use --max-id)")
|
||||
|
||||
if os.path.exists(args.out):
|
||||
os.remove(args.out)
|
||||
cur.execute("ATTACH DATABASE ? AS inc", (args.out,))
|
||||
|
||||
where = "id>?"
|
||||
params = [args.id_start]
|
||||
if args.id_end is not None:
|
||||
where += " AND id<=?"
|
||||
params.append(args.id_end)
|
||||
|
||||
cur.execute(
|
||||
"CREATE TABLE inc.dbbardata AS "
|
||||
"SELECT %s FROM dbbardata WHERE %s" % (COLS, where),
|
||||
params)
|
||||
n = cur.execute("SELECT COUNT(*) FROM inc.dbbardata").fetchone()[0]
|
||||
# 导出范围的 max id(从源库查,rowid 快),供 NAS 记为下次起点
|
||||
max_exported = cur.execute(
|
||||
"SELECT MAX(id) FROM dbbardata WHERE %s" % where, params).fetchone()[0]
|
||||
conn.commit()
|
||||
conn.close()
|
||||
size_mb = os.path.getsize(args.out) / 1048576.0
|
||||
print("EXPORTED rows=%d max_id=%d sizeMB=%.1f -> %s"
|
||||
% (n, max_exported if max_exported else 0, size_mb, args.out))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,55 @@
|
||||
"""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()
|
||||
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()
|
||||
@@ -0,0 +1,91 @@
|
||||
#!/bin/bash
|
||||
# NAS 端:VPS dbbardata 增量/全量同步(id 增量,rowid 索引快;VPS pull 方向)。
|
||||
#
|
||||
# 用法:
|
||||
# 首次全量(分片,断点续传): bash sync_dbbardata.sh full
|
||||
# 每日增量(定时任务): bash sync_dbbardata.sh increment
|
||||
#
|
||||
# 设计:
|
||||
# - VPS 端脚本不落盘,通过 ssh stdin 喂 export_increment.py 执行(绕 Windows 部署)。
|
||||
# - 全量按 id 分片(STEP 行/片)逐片导出→pull→merge,since.txt 记进度,中断重跑 full 自动续。
|
||||
# - 增量:since=since.txt 的 max_id,导出 id>since,更新 since.txt。
|
||||
# - merge 用 INSERT OR IGNORE,UNIQUE(symbol,exchange,interval,datetime) 去重,重跑安全。
|
||||
set -eu
|
||||
|
||||
KEY=/var/services/homes/admin/.ssh/id_ed25519_nas
|
||||
VPS=Administrator@49.232.102.198
|
||||
VPS_DB='C:\sanguo_vnpy_v2\data\quant_trading.db'
|
||||
VPS_OUT='C:\sanguo_vnpy_v2\data\_sync_inc.db'
|
||||
VPS_OUT_SCP='C:/sanguo_vnpy_v2/data/_sync_inc.db'
|
||||
VPS_PY='C:\Python310\python.exe -X utf8 -'
|
||||
|
||||
ROOT=/volume1/stock/sanguo_vnpy_v2
|
||||
EXP=$ROOT/scripts/nas_sync/export_increment.py
|
||||
MERGE=$ROOT/scripts/nas_sync/merge_increment.py
|
||||
DB=$ROOT/data_backup/quant_trading.db
|
||||
STAGE=$ROOT/data_backup/_staging
|
||||
STATE=$ROOT/data_backup/since.txt
|
||||
LOG=$ROOT/data_backup/sync.log
|
||||
STEP=20000000 # 全量分片步长(行/片,~2GB/片)
|
||||
|
||||
mkdir -p "$STAGE"
|
||||
LOG_TS="$(date '+%F %T')"
|
||||
echo "=== $LOG_TS MODE=${1:-increment} ===" >> "$LOG"
|
||||
|
||||
get_max() {
|
||||
ssh -i "$KEY" -o StrictHostKeyChecking=no "$VPS" "$VPS_PY --db $VPS_DB --max-id" < "$EXP" \
|
||||
| grep -oE 'MAX_ID=[0-9]+' | cut -d= -f2
|
||||
}
|
||||
|
||||
export_slice() {
|
||||
# $1=id_start $2=id_end(空=不限)
|
||||
if [ -n "$2" ]; then
|
||||
ssh -i "$KEY" -o StrictHostKeyChecking=no "$VPS" \
|
||||
"$VPS_PY --db $VPS_DB --out $VPS_OUT --id-start $1 --id-end $2" < "$EXP"
|
||||
else
|
||||
ssh -i "$KEY" -o StrictHostKeyChecking=no "$VPS" \
|
||||
"$VPS_PY --db $VPS_DB --out $VPS_OUT --id-start $1" < "$EXP"
|
||||
fi
|
||||
}
|
||||
|
||||
pull_merge() {
|
||||
scp -i "$KEY" -o StrictHostKeyChecking=no "$VPS:$VPS_OUT_SCP" "$STAGE/inc.db" >> "$LOG" 2>&1
|
||||
python3 "$MERGE" --db "$DB" --inc "$STAGE/inc.db" >> "$LOG" 2>&1
|
||||
rm -f "$STAGE/inc.db"
|
||||
}
|
||||
|
||||
CUR_MAX=$(get_max)
|
||||
echo "[$LOG_TS] VPS max_id=$CUR_MAX" >> "$LOG"
|
||||
|
||||
case "${1:-increment}" in
|
||||
full)
|
||||
SINCE=$(cat "$STATE" 2>/dev/null || echo 0)
|
||||
echo "[$LOG_TS] full since=$SINCE cur_max=$CUR_MAX step=$STEP" >> "$LOG"
|
||||
while [ "$SINCE" -lt "$CUR_MAX" ]; do
|
||||
END=$((SINCE + STEP))
|
||||
[ "$END" -gt "$CUR_MAX" ] && END=$CUR_MAX
|
||||
echo "[$(date '+%T')] slice id($SINCE,$END]" >> "$LOG"
|
||||
export_slice "$SINCE" "$END"
|
||||
pull_merge
|
||||
echo "$END" > "$STATE"
|
||||
SINCE=$END
|
||||
done
|
||||
echo "[$(date '+%T')] FULL DONE since=$SINCE" >> "$LOG"
|
||||
;;
|
||||
increment)
|
||||
SINCE=$(cat "$STATE" 2>/dev/null || echo 0)
|
||||
echo "[$LOG_TS] increment since=$SINCE -> $CUR_MAX" >> "$LOG"
|
||||
if [ "$SINCE" -ge "$CUR_MAX" ]; then
|
||||
echo "[$LOG_TS] up-to-date" >> "$LOG"
|
||||
exit 0
|
||||
fi
|
||||
export_slice "$SINCE" ""
|
||||
pull_merge
|
||||
echo "$CUR_MAX" > "$STATE"
|
||||
echo "[$(date '+%T')] INCREMENT DONE since=$CUR_MAX" >> "$LOG"
|
||||
;;
|
||||
*)
|
||||
echo "usage: $0 {full|increment}" >&2
|
||||
exit 1
|
||||
;;
|
||||
esac
|
||||
Reference in New Issue
Block a user