Compare commits
2 Commits
6aec9edf81
...
b2aab85105
| Author | SHA1 | Date | |
|---|---|---|---|
| b2aab85105 | |||
| 176c17f611 |
@@ -7,6 +7,7 @@ POST /portfolio/backtest: SSH 触发 VPS 跑 BulletTrade + all_weather,
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import subprocess
|
||||
@@ -54,18 +55,29 @@ def run_portfolio_backtest(req: PortfolioBacktestRequest):
|
||||
"""
|
||||
# 本地模式(VPS 后端直接跑 runner,避免 SSH 自连外网 IP 绕路);
|
||||
# SSH 模式(Mac 后端 → VPS)。env SANGUO_PORTFOLIO_LOCAL=1 切本地。
|
||||
# 自动检测:VPS 上 _VPS_WORKDIR(C:\sanguo_vnpy_v2) 存在 → 本地直接跑 runner;
|
||||
# Mac 上无该目录 → SSH 到 VPS。无需 env 配置,同源代码两端自适应。
|
||||
if os.path.isdir(_VPS_WORKDIR):
|
||||
# 自动检测三机自适应:
|
||||
# - VPS(Windows): _VPS_WORKDIR(C:\sanguo_vnpy_v2) 存在 → 本地直接跑 runner(默认 provider,行为不变)
|
||||
# - NAS(Linux 容器): /app 存在(docker-compose 代码挂载目录) → 本地跑 runner --provider unified 读 NAS 权威数据层
|
||||
# - Mac(dev): 都不满足 → SSH 到 VPS(原行为)
|
||||
if os.path.isdir(_VPS_WORKDIR) or os.path.isdir("/app"):
|
||||
local_cwd = _VPS_WORKDIR if os.path.isdir(_VPS_WORKDIR) else "/app"
|
||||
argv = [
|
||||
sys.executable, "-X", "utf8", "-m", "sanguo_portfolio.runner_backtest",
|
||||
"--json", "--start", req.start_date, "--end", req.end_date,
|
||||
"--cash", str(req.initial_cash), "--benchmark", req.benchmark, "--max-pool", "30",
|
||||
]
|
||||
# NAS 容器分支: unified provider 读 NAS 权威数据层(dbbardata + parquet);
|
||||
# VPS 本地分支保持原状(默认 provider,cwd=_VPS_WORKDIR,行为零变化)
|
||||
if local_cwd == "/app":
|
||||
nas_provider_config = json.dumps({
|
||||
"db_path": "/volume1/stock/sanguo_vnpy_v2/data_backup/quant_trading.db",
|
||||
"data_dir": "/volume1/stock/sanguo_vnpy_v2/data",
|
||||
})
|
||||
argv += ["--provider", "unified", "--provider-config", nas_provider_config]
|
||||
logger.info("[portfolio] 本地跑 runner: %s", " ".join(argv[3:]))
|
||||
try:
|
||||
proc = subprocess.run(
|
||||
argv, cwd=_VPS_WORKDIR, capture_output=True, text=True,
|
||||
argv, cwd=local_cwd, capture_output=True, text=True,
|
||||
timeout=_VPS_TIMEOUT, check=False,
|
||||
)
|
||||
except subprocess.TimeoutExpired:
|
||||
|
||||
@@ -344,7 +344,7 @@ def run_backtest_json(params: Dict[str, Any]) -> Dict[str, Any]:
|
||||
frequency="day",
|
||||
strategy=strategy_name,
|
||||
provider=params.get("provider", "local"),
|
||||
provider_config="{}",
|
||||
provider_config=params.get("provider_config", "{}"),
|
||||
result_file="", # JSON 模式不写 md
|
||||
max_pool=int(params.get("max_pool", 0)),
|
||||
)
|
||||
@@ -503,6 +503,7 @@ def main() -> None:
|
||||
"initial_cash": args.cash,
|
||||
"benchmark": args.benchmark,
|
||||
"provider": args.provider,
|
||||
"provider_config": args.provider_config,
|
||||
"max_pool": args.max_pool,
|
||||
})
|
||||
print(json.dumps(result, ensure_ascii=False, default=str))
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
"""VPS 端:全量导出单张小表到独立 sqlite 文件(搬数据,merge 端重建 schema)。
|
||||
|
||||
用于 constituent_unified / bs_adjust_factor 等小表全量同步
|
||||
(无 id 增量键、行数万级、变化不频繁,全量 dump+replace 最简)。
|
||||
|
||||
- 纯 sqlite3 标准库,零第三方依赖。
|
||||
- 只读主库(mode=ro),不干扰 VPS 生产写入。
|
||||
- CREATE TABLE AS SELECT 搬数据(列名保留;约束由 NAS merge 端权威 schema 重建)。
|
||||
|
||||
用法:
|
||||
python export_table.py --db PATH --out out.db --table constituent_unified
|
||||
"""
|
||||
import argparse
|
||||
import os
|
||||
import sqlite3
|
||||
|
||||
|
||||
def main():
|
||||
ap = argparse.ArgumentParser(description="全量导出单张小表到独立 sqlite")
|
||||
ap.add_argument("--db", required=True, help="源 quant_trading.db 路径")
|
||||
ap.add_argument("--out", required=True, help="输出 sqlite 路径")
|
||||
ap.add_argument("--table", required=True, help="表名")
|
||||
args = ap.parse_args()
|
||||
tbl = args.table
|
||||
|
||||
conn = sqlite3.connect("file:%s?mode=ro" % args.db, uri=True)
|
||||
cur = conn.cursor()
|
||||
exists = cur.execute(
|
||||
"SELECT 1 FROM sqlite_master WHERE type='table' AND name=?", (tbl,)
|
||||
).fetchone()
|
||||
if not exists:
|
||||
raise SystemExit("table not found in source db: %s" % tbl)
|
||||
|
||||
if os.path.exists(args.out):
|
||||
os.remove(args.out)
|
||||
cur.execute("ATTACH DATABASE ? AS inc", (args.out,))
|
||||
cur.execute(
|
||||
"CREATE TABLE inc.[%s] AS SELECT * FROM main.[%s]" % (tbl, tbl)
|
||||
)
|
||||
n = cur.execute("SELECT COUNT(*) FROM inc.[%s]" % tbl).fetchone()[0]
|
||||
conn.commit()
|
||||
conn.close()
|
||||
print("EXPORTED table=%s rows=%d sizeMB=%.2f -> %s"
|
||||
% (tbl, n, os.path.getsize(args.out) / 1048576.0, args.out))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,79 @@
|
||||
"""NAS 端:全量替换单张小表(DROP + CREATE + INSERT),幂等可重跑。
|
||||
|
||||
用于 constituent_unified / bs_adjust_factor 等小表。
|
||||
- 纯 sqlite3 标准库,零第三方依赖(NAS 宿主 Python 3.8 也能跑)。
|
||||
- DROP IF EXISTS + 权威 schema 重建 + INSERT:每次全量替换,重跑安全。
|
||||
- 权威 schema(含主键/约束)硬编码在此,NAS 副本对齐 VPS。
|
||||
|
||||
用法:
|
||||
python merge_table.py --db /volume1/.../quant_trading.db --inc /tmp/t.db --table constituent_unified
|
||||
"""
|
||||
import argparse
|
||||
import sqlite3
|
||||
|
||||
# 权威 schema(对齐 VPS quant_trading.db)。新增表在此登记。
|
||||
SCHEMAS = {
|
||||
"constituent_unified": (
|
||||
"CREATE TABLE constituent_unified("
|
||||
"index_code TEXT, code TEXT, code_name TEXT, source TEXT,"
|
||||
"in_current INT, was_removed INT)"
|
||||
),
|
||||
"bs_adjust_factor": (
|
||||
"CREATE TABLE bs_adjust_factor("
|
||||
"code TEXT NOT NULL, dividOperateDate TEXT NOT NULL,"
|
||||
"foreAdjustFactor REAL, backAdjustFactor REAL, adjustFactor REAL,"
|
||||
"PRIMARY KEY (code, dividOperateDate))"
|
||||
),
|
||||
}
|
||||
# 列名(与源表对齐,用于 INSERT 列顺序)。
|
||||
COLUMNS = {
|
||||
"constituent_unified":
|
||||
"index_code,code,code_name,source,in_current,was_removed",
|
||||
"bs_adjust_factor":
|
||||
"code,dividOperateDate,foreAdjustFactor,backAdjustFactor,adjustFactor",
|
||||
}
|
||||
# 额外索引(主键自带的不列)。
|
||||
INDEXES = {
|
||||
"constituent_unified": [
|
||||
"CREATE INDEX IF NOT EXISTS idx_cu_index ON constituent_unified(index_code)",
|
||||
"CREATE INDEX IF NOT EXISTS idx_cu_code ON constituent_unified(code)",
|
||||
],
|
||||
"bs_adjust_factor": [],
|
||||
}
|
||||
|
||||
|
||||
def main():
|
||||
ap = argparse.ArgumentParser(description="全量替换单张小表进 NAS 副本")
|
||||
ap.add_argument("--db", required=True, help="NAS quant_trading.db 路径")
|
||||
ap.add_argument("--inc", required=True, help="增量 sqlite 路径")
|
||||
ap.add_argument("--table", required=True, help="表名")
|
||||
args = ap.parse_args()
|
||||
tbl = args.table
|
||||
if tbl not in SCHEMAS:
|
||||
raise SystemExit("unknown table %s; add schema/columns first" % tbl)
|
||||
cols = COLUMNS[tbl]
|
||||
|
||||
conn = sqlite3.connect(args.db)
|
||||
cur = conn.cursor()
|
||||
for p in ("PRAGMA journal_mode=WAL", "PRAGMA synchronous=NORMAL"):
|
||||
cur.execute(p)
|
||||
cur.execute("ATTACH DATABASE ? AS inc", (args.inc,))
|
||||
|
||||
inc_count = cur.execute(
|
||||
"SELECT COUNT(*) FROM inc.[%s]" % tbl).fetchone()[0]
|
||||
cur.execute("DROP TABLE IF EXISTS main.[%s]" % tbl)
|
||||
cur.execute(SCHEMAS[tbl])
|
||||
for idx in INDEXES.get(tbl, []):
|
||||
cur.execute(idx)
|
||||
cur.execute(
|
||||
"INSERT INTO main.[%s](%s) SELECT %s FROM inc.[%s]" % (tbl, cols, cols, tbl)
|
||||
)
|
||||
conn.commit()
|
||||
after = cur.execute(
|
||||
"SELECT COUNT(*) FROM main.[%s]" % tbl).fetchone()[0]
|
||||
conn.close()
|
||||
print("REPLACED table=%s inc=%d after=%d" % (tbl, inc_count, after))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -89,3 +89,9 @@ case "${1:-increment}" in
|
||||
exit 1
|
||||
;;
|
||||
esac
|
||||
|
||||
# 小表全量同步(constituent_unified 成份股 / bs_adjust_factor 复权因子:
|
||||
# 月度或除权日才变,万级行,全量几秒)。独立 log(sync_tables.log),
|
||||
# 失败不影响 dbbardata 主流程。每日 increment 连带跑,保持 NAS 副本不陈旧。
|
||||
bash "$(dirname "$0")/sync_tables.sh" >> "$LOG" 2>&1 \
|
||||
|| echo "[$(date '+%T')] WARN sync_tables 非致命" >> "$LOG"
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
#!/bin/bash
|
||||
# NAS 端:小表全量同步(constituent_unified / bs_adjust_factor 等)。
|
||||
#
|
||||
# 这些表无 id 增量键、行数万级、变化不频繁(成份股月度调整 / 除权日),
|
||||
# 全量 dump + replace 最简,每天跑也只需几秒。
|
||||
#
|
||||
# 用法: bash sync_tables.sh
|
||||
# 设计同 sync_dbbardata.sh: VPS python(stdin 喂) 导出 sqlite → scp → NAS 全量替换。
|
||||
# 幂等: merge 端 DROP+CREATE+INSERT,重跑安全。
|
||||
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_tbl.db'
|
||||
VPS_OUT_SCP='C:/sanguo_vnpy_v2/data/_sync_tbl.db'
|
||||
VPS_PY='C:\Python310\python.exe -X utf8 -'
|
||||
|
||||
ROOT=/volume1/stock/sanguo_vnpy_v2
|
||||
EXP=$ROOT/scripts/nas_sync/export_table.py
|
||||
MERGE=$ROOT/scripts/nas_sync/merge_table.py
|
||||
DB=$ROOT/data_backup/quant_trading.db
|
||||
STAGE=$ROOT/data_backup/_staging
|
||||
LOG=$ROOT/data_backup/sync_tables.log
|
||||
|
||||
# 要全量同步的小表清单(merge_table.py 需登记 schema)
|
||||
TABLES=(constituent_unified bs_adjust_factor)
|
||||
|
||||
mkdir -p "$STAGE"
|
||||
echo "=== $(date '+%F %T') tables sync start ===" >> "$LOG"
|
||||
|
||||
for T in "${TABLES[@]}"; do
|
||||
echo "[$(date '+%T')] export $T ..." >> "$LOG"
|
||||
ssh -i "$KEY" -o StrictHostKeyChecking=no "$VPS" \
|
||||
"$VPS_PY --db $VPS_DB --out $VPS_OUT --table $T" < "$EXP" >> "$LOG" 2>&1
|
||||
scp -i "$KEY" -o StrictHostKeyChecking=no "$VPS:$VPS_OUT_SCP" "$STAGE/tbl.db" >> "$LOG" 2>&1
|
||||
python3 "$MERGE" --db "$DB" --inc "$STAGE/tbl.db" --table "$T" >> "$LOG" 2>&1
|
||||
rm -f "$STAGE/tbl.db"
|
||||
echo "[$(date '+%T')] $T done" >> "$LOG"
|
||||
done
|
||||
echo "=== $(date '+%F %T') TABLES SYNC DONE ===" >> "$LOG"
|
||||
Reference in New Issue
Block a user