feat(data): 全市场5m/15m下载+灌库+每日增量链

xtdata逐只下载(非批量download_history_data2,后者5201只触发xtquant死锁),
先5m后15m(15m依赖5m),看get bars判成败(ret=None正常非失败)。
download_15m_xtdata.py + relaunch_15m_wrapper.ps1(auto-restart兜底segfault)。
import_vnpy_minute_fast.py灌库(INSERT OR REPLACE,interval存5m/15m,8028万行)。
_run_daily.ps1加5m/15m增量段(每天16:30)。validate_import.py校验。

全市场5201只,5m 6020万/15m 2007万bar,2025-07-17~2026-07-17。
This commit is contained in:
2026-07-18 18:45:24 +08:00
parent 43b6a58ba7
commit ee27313644
8 changed files with 1472 additions and 0 deletions
+6
View File
@@ -0,0 +1,6 @@
@echo off
REM Launch download_15m_xtdata.py detached. Logs to data\minute_15\download.log
cd /d C:\sanguo_vnpy_v2\scripts\data_platform
set PYTHONIOENCODING=utf-8
C:\Python310\python.exe -X utf8 download_15m_xtdata.py > C:\sanguo_vnpy_v2\data\minute_15\download.log 2>&1
exit /b %ERRORLEVEL%
+9
View File
@@ -0,0 +1,9 @@
# schtasks /ru SYSTEM 下跑 build_daily_from_xtdata.py 的 wrapper。
# 输出 tee 到日志,便于 ssh 断开后看进度。不依赖工作目录(脚本内 ROOT 硬编码)。
$ErrorActionPreference = 'Continue'
$log = 'C:\sanguo_vnpy_v2\data\xtdata_build.log'
$py = 'C:\Python310\python.exe'
$script = 'C:\sanguo_vnpy_v2\scripts\data_platform\build_daily_from_xtdata.py'
"=== BUILD START $(Get-Date -Format o) ===" | Out-File -FilePath $log -Encoding utf8
& $py -X utf8 $script *>&1 | Tee-Object -FilePath $log -Append
"=== BUILD END $(Get-Date -Format o) exit=$LASTEXITCODE ===" | Out-File -FilePath $log -Encoding utf8 -Append
+272
View File
@@ -0,0 +1,272 @@
# schtasks 每日定时跑:下载日线 → 灌库 → 校验(串联,按退出码守卫)。收盘后触发(如 16:30)。
# 下载失败绝不灌库;灌库失败绝不校验;校验失败 schtask LastTaskResult 反映为 1。
$ErrorActionPreference = 'Continue'
$log = 'C:\sanguo_vnpy_v2\data\daily_update.log'
$import_log = 'C:\sanguo_vnpy_v2\data\import_db.log'
$validate_log = 'C:\sanguo_vnpy_v2\data\validate.log'
$minute_log = 'C:\sanguo_vnpy_v2\data\minute_update.log'
$minute_import_log = 'C:\sanguo_vnpy_v2\data\minute_import.log'
$minute_validate_log = 'C:\sanguo_vnpy_v2\data\minute_validate.log'
$py = 'C:\Python310\python.exe'
$download_script = 'C:\sanguo_vnpy_v2\scripts\data_platform\daily_update_xtdata.py'
$import_script = 'C:\sanguo_vnpy_v2\scripts\data_platform\import_vnpy_daily_fast.py'
$validate_script = 'C:\sanguo_vnpy_v2\scripts\data_platform\validate_import.py'
$minute_import_script = 'C:\sanguo_vnpy_v2\scripts\data_platform\import_vnpy_minute_fast.py'
# 全市场 5m/15m 增量下载脚本(inlinextquant 单线程 paced,先 5m 后 15m
# 只下最近 7 天,已存在且 max(datetime)>=today 则跳过该 code。
$minute_incr_script = @"
import os, sys, time, datetime as dt
import pandas as pd
from xtquant import xtdata as xd
ROOT = r'C:\sanguo_vnpy_v2'
LOOKBACK_DAYS = 7
END = dt.datetime.now().strftime('%Y%m%d')
START = (dt.datetime.now() - dt.timedelta(days=LOOKBACK_DAYS)).strftime('%Y%m%d')
TODAY = dt.datetime.now().strftime('%Y-%m-%d')
RAW_5M = os.path.join(ROOT, 'data', 'minute_5', 'raw')
RAW_15M = os.path.join(ROOT, 'data', 'minute_15', 'raw')
SLEEP_5M = 0.05
SLEEP_15M = 0.2
PROGRESS_EVERY = 200
def log(m):
print(f'[INCR-MIN {time.time():.0f}] {m}', flush=True)
def parse_dt_index(idx):
s = pd.Series([str(i) for i in idx]).str.slice(0, 14)
return pd.to_datetime(s, format='%Y%m%d%H%M%S', errors='coerce')
def to_pdf(df):
return pd.DataFrame({
'datetime': parse_dt_index(df.index),
'open': df['open'].astype(float).values,
'high': df['high'].astype(float).values,
'low': df['low'].astype(float).values,
'close': df['close'].astype(float).values,
'volume': (df['volume'].astype(float) * 100.0).values,
}).dropna(subset=['datetime']).sort_values('datetime').reset_index(drop=True)
def atomic_write(path, df):
os.makedirs(os.path.dirname(path), exist_ok=True)
tmp = path + '.tmp'
df.to_parquet(tmp, index=False)
os.replace(tmp, path)
def fetch(code, period):
r = xd.get_market_data_ex([], [code], period=period, start_time=START,
end_time=END, dividend_type='none')
return r.get(code) if r else None
def process_one(code, period, raw_dir):
sym, exc = code.split('.')
raw_path = os.path.join(raw_dir, f'{sym}.{exc}_{period}.parquet')
# max(datetime)>=today
if os.path.exists(raw_path):
try:
old = pd.read_parquet(raw_path, columns=['datetime'])
if len(old):
last_dt = str(old['datetime'].max())[:10]
if last_dt >= TODAY:
return 'skipped', len(old)
except Exception:
pass
# downloadret=None get bars
try:
xd.download_history_data(code, period, START, END)
except Exception as e:
return f'dl_err:{e}', 0
# get raw merge / 7
try:
rdf = fetch(code, period)
if rdf is None or not len(rdf):
return 'no_data', 0
pdf = to_pdf(rdf)
if not len(pdf):
return 'empty_pdf', 0
if os.path.exists(raw_path):
try:
old = pd.read_parquet(raw_path)
merged = (pd.concat([old, pdf])
.drop_duplicates('datetime', keep='last')
.sort_values('datetime')
.reset_index(drop=True))
except Exception:
merged = pdf
else:
merged = pdf
atomic_write(raw_path, merged)
return 'ok', len(merged)
except Exception as e:
return f'write_err:{e}', 0
def phase(period, raw_dir, sleep_s, universe):
log(f'PHASE {period} START n={len(universe)} window={START}~{END}')
ok = skip = fail = 0
for i, code in enumerate(universe):
try:
status, n = process_one(code, period, raw_dir)
if status == 'ok':
ok += 1
elif status == 'skipped':
skip += 1
else:
fail += 1
if fail <= 5:
log(f' {code} {period} {status}')
except Exception as e:
fail += 1
if fail <= 5:
log(f' {code} {period} exc: {e}')
if (i + 1) % PROGRESS_EVERY == 0 or (i + 1) == len(universe):
log(f' {period} {i+1}/{len(universe)} ok={ok} skip={skip} fail={fail}')
time.sleep(sleep_s)
log(f'PHASE {period} DONE ok={ok} skip={skip} fail={fail}')
return ok, skip, fail
def main():
log(f'START window={START}~{END} today={TODAY}')
try:
u = xd.get_stock_list_in_sector('A') or []
except Exception as e:
log(f'FATAL get_stock_list: {e}')
os._exit(2)
if not u:
log('FATAL: empty universe (miniQMT ?)')
os._exit(2)
log(f'universe={len(u)}')
ok5, skip5, fail5 = phase('5m', RAW_5M, SLEEP_5M, u)
ok15, skip15, fail15 = phase('15m', RAW_15M, SLEEP_15M, u)
total_fail = fail5 + fail15
log(f'INCR DONE 5m(ok={ok5},skip={skip5},fail={fail5}) '
f'15m(ok={ok15},skip={skip15},fail={fail15}) total_fail={total_fail}')
sys.stdout.flush()
# /退退 1
os._exit(0 if total_fail < len(u) else 1)
main()
"@
# 全市场 5m/15m 增量校验脚本(inline
$minute_validate_script = @"
import os, sys, datetime as dt, sqlite3
import pandas as pd
DB = r'C:\sanguo_vnpy_v2\data\quant_trading.db'
TODAY = dt.datetime.now().strftime('%Y-%m-%d')
SAMPLES = ['600519', '000001', '000858']
EXPECT_INTERVALS = ['5m', '15m']
def log(m):
print(f'[VAL-MIN] {m}', flush=True)
def main():
if not os.path.exists(DB):
log(f'FAIL: DB not found {DB}')
sys.exit(1)
conn = sqlite3.connect(DB)
today = TODAY
log(f'today={today}')
ok = True
for iv in EXPECT_INTERVALS:
row = conn.execute(
'SELECT MAX(datetime), COUNT(*) FROM dbbardata WHERE interval=?', (iv,)
).fetchone()
max_dt, total = row
max_date = (max_dt or '')[:10]
iv_ok = (max_date == today) and (total or 0) > 0
log(f'[{iv}] total={total} max_dt={max_dt} max_date={max_date} '
f'today_match={max_date == today} -> {"PASS" if iv_ok else "FAIL"}')
if not iv_ok:
ok = False
for sym in SAMPLES:
row = conn.execute(
"SELECT COUNT(*) FROM dbbardata WHERE symbol=? AND interval='5m' "
"AND datetime LIKE ?", (sym, today + '%')
).fetchone()
n = row[0] or 0
log(f'sample {sym} 5m today_bars={n} -> {"PASS" if n > 0 else "FAIL"}')
if n <= 0:
ok = False
conn.close()
log(f'OVERALL {"PASS" if ok else "FAIL"}')
sys.exit(0 if ok else 1)
main()
"@
"=== DAILY UPDATE START $(Get-Date -Format o) ===" | Out-File -FilePath $log -Encoding utf8
& $py -X utf8 $download_script *>&1 | Tee-Object -FilePath $log -Append
$download_exit = $LASTEXITCODE
"=== DAILY UPDATE END $(Get-Date -Format o) exit=$download_exit ===" | Out-File -FilePath $log -Encoding utf8 -Append
if ($download_exit -ne 0) {
"=== IMPORT SKIPPED (download failed) $(Get-Date -Format o) ===" | Out-File -FilePath $log -Encoding utf8 -Append
exit 1
}
# 灌库(全量 2010 起,幂等 INSERT OR REPLACE,重跑安全)
$env:VNPY_DB_PATH = 'C:\sanguo_vnpy_v2\data\quant_trading.db'
$env:DAILY_DIR = 'C:\sanguo_vnpy_v2\data\raw'
"=== IMPORT START $(Get-Date -Format o) ===" | Out-File -FilePath $import_log -Encoding utf8
& $py -X utf8 $import_script --start-year 2010 *>&1 | Tee-Object -FilePath $import_log -Append
$import_exit = $LASTEXITCODE
"=== IMPORT DONE $(Get-Date -Format o) exit=$import_exit ===" | Out-File -FilePath $import_log -Encoding utf8 -Append
"=== IMPORT DONE exit=$import_exit ===" | Out-File -FilePath $log -Encoding utf8 -Append
if ($import_exit -ne 0) {
"=== VALIDATE SKIPPED (import failed) $(Get-Date -Format o) ===" | Out-File -FilePath $log -Encoding utf8 -Append
exit 1
}
# 灌后校验
"=== VALIDATE START $(Get-Date -Format o) ===" | Out-File -FilePath $validate_log -Encoding utf8
& $py -X utf8 $validate_script *>&1 | Tee-Object -FilePath $validate_log -Append
$validate_exit = $LASTEXITCODE
"=== VALIDATE DONE $(Get-Date -Format o) exit=$validate_exit ===" | Out-File -FilePath $validate_log -Encoding utf8 -Append
"=== VALIDATE DONE exit=$validate_exit ===" | Out-File -FilePath $log -Encoding utf8 -Append
# ===================== 5m / 15m 增量下载 → 灌库 → 校验 =====================
# 守卫:日线 validate 失败仍跑 minute(独立链路),但 daily 失败码在最终 exit 中体现
$env:VNPY_DB_PATH = 'C:\sanguo_vnpy_v2\data\quant_trading.db'
$minute_final_exit = 0
"=== MINUTE DOWNLOAD START $(Get-Date -Format o) ===" | Out-File -FilePath $minute_log -Encoding utf8
"DiskFree_GB=$((Get-PSDrive C).Free / 1GB)" | Out-File -FilePath $minute_log -Encoding utf8 -Append
# inline python via stdin(约束:不新建文件,inline 在 ps1 内)
$minute_incr_script | & $py -X utf8 - *>&1 | Tee-Object -FilePath $minute_log -Append
$minute_download_exit = $LASTEXITCODE
"=== MINUTE DOWNLOAD END $(Get-Date -Format o) exit=$minute_download_exit ===" | Out-File -FilePath $minute_log -Encoding utf8 -Append
"=== MINUTE DOWNLOAD DONE exit=$minute_download_exit ===" | Out-File -FilePath $log -Encoding utf8 -Append
if ($minute_download_exit -ne 0) {
"=== MINUTE IMPORT SKIPPED (download failed) $(Get-Date -Format o) ===" | Out-File -FilePath $minute_log -Encoding utf8 -Append
$minute_final_exit = $minute_download_exit
} else {
# 灌库(全量 upsert 已存在+新增,幂等)
"=== MINUTE IMPORT START $(Get-Date -Format o) ===" | Out-File -FilePath $minute_import_log -Encoding utf8
& $py -X utf8 $minute_import_script *>&1 | Tee-Object -FilePath $minute_import_log -Append
$minute_import_exit = $LASTEXITCODE
"=== MINUTE IMPORT DONE $(Get-Date -Format o) exit=$minute_import_exit ===" | Out-File -FilePath $minute_import_log -Encoding utf8 -Append
"=== MINUTE IMPORT DONE exit=$minute_import_exit ===" | Out-File -FilePath $log -Encoding utf8 -Append
if ($minute_import_exit -ne 0) {
"=== MINUTE VALIDATE SKIPPED (import failed) $(Get-Date -Format o) ===" | Out-File -FilePath $minute_log -Encoding utf8 -Append
$minute_final_exit = $minute_import_exit
} else {
# 灌后校验:5m/15m max(datetime)==today + 抽样 bars>0
"=== MINUTE VALIDATE START $(Get-Date -Format o) ===" | Out-File -FilePath $minute_validate_log -Encoding utf8
$minute_validate_script | & $py -X utf8 - *>&1 | Tee-Object -FilePath $minute_validate_log -Append
$minute_validate_exit = $LASTEXITCODE
"=== MINUTE VALIDATE DONE $(Get-Date -Format o) exit=$minute_validate_exit ===" | Out-File -FilePath $minute_validate_log -Encoding utf8 -Append
"=== MINUTE VALIDATE DONE exit=$minute_validate_exit ===" | Out-File -FilePath $log -Encoding utf8 -Append
$minute_final_exit = $minute_validate_exit
}
}
"DiskFree_GB_end=$((Get-PSDrive C).Free / 1GB)" | Out-File -FilePath $log -Encoding utf8 -Append
# 总 exit:daily 与 minute 任一非零则非零(对应 schtasks LastTaskResult 反映状态)
$final_exit = $validate_exit
if ($minute_final_exit -ne 0) { $final_exit = $minute_final_exit }
exit $final_exit
+6
View File
@@ -0,0 +1,6 @@
$env:VNPY_DB_PATH = "C:\sanguo_vnpy_v2\data\quant_trading.db"
$env:DAILY_DIR = "C:\sanguo_vnpy_v2\data\raw"
$log = "C:\sanguo_vnpy_v2\data\import_db.log"
"=== IMPORT START $(Get-Date -Format o) ===" | Out-File -FilePath $log -Encoding utf8
& "C:\Python310\python.exe" -X utf8 "C:\sanguo_vnpy_v2\scripts\data_platform\import_vnpy_daily_fast.py" --start-year 2010 *>&1 | Tee-Object -FilePath $log -Append
"=== IMPORT DONE $(Get-Date -Format o) exit=$LASTEXITCODE ===" | Out-File -FilePath $log -Encoding utf8 -Append
@@ -0,0 +1,637 @@
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""全市场 5m + 15m 双周期一次性下载(raw + qfq 双源)。VPS 后台跑(需 miniQMT 常驻)。
【Main Agent 实证真根因 — 勿再质疑】
1. download_history_data / download_history_data2 返回 None 是正常的,绝非失败!
成败只看下一步 get_market_data_ex 的 bars 数(>0=成功)。
2. 15m 依赖 5m 基础数据:先下 5m 再下 15m(顺序不可反、不可并行)。
实证:先 5m=11670 bars、再 15m=3890 bars(全板块一致)。
3. 全市场 1m/5m 直接可下(不需订阅、不是权限问题)。
【复权机制】download 只下原始 none 数据一次,读时用 dividend_type='front' 转 qfq。
每周期 download 一次,raw/qfq 两份在读时分流出。别为 qfq 单独 download。
(同 daily_update_xtdata.pydownload_history_data2 批量 → get_market_data_ex 分别
读 none/front 两份)。
【两轮下载顺序】
Phase 1: 全市场 download_history_data2 批量下 5m(→ xtdata 本地缓存)
Phase 1.5: 逐只 get_market_data_ex 读 5mraw + qfq)→ 写 parquet
Phase 2: 全市场 download_history_data2 批量下 15m(依赖 5m 基础)
Phase 2.5: 逐只 get_market_data_ex 读 15mraw + qfq)→ 写 parquet
Phase 3: 校验 → 写 _result.json
【存储】C:\\sanguo_vnpy_v2\\data\\minute_5\\{raw,qfq}\\<code>_5m.parquet
C:\\sanguo_vnpy_v2\\data\\minute_15\\{raw,qfq}\\<code>_15m.parquet
schema: datetime,open,high,low,close,volumevolume ×100 手→股)。
【断点续传】文件已存在且 bars ≥ MIN_BARS_RESUME_{5M,15M} 则跳过该只写盘。
【单线程 paced】批量 download_history_data2 服务端内部并行;轮次间 sleep 2s。
【进度】_progress.json 每 200 只刷新。
用法(VPS):
C:\\Python310\\python.exe -X utf8 download_15m_xtdata.py
ETF 本轮未下(universe 取"沪深A股"约 5201 只)。
退出码:0=PASS;1=FAIL 或有失败;2=致命错误(universe 空 / miniQMT 未连)。
"""
import os
import sys
import json
import time
import random
import datetime as dt
from typing import Any
import pandas as pd
ROOT = r"C:\sanguo_vnpy_v2"
DATA_5M = os.path.join(ROOT, "data", "minute_5")
DATA_15M = os.path.join(ROOT, "data", "minute_15")
RAW_5M = os.path.join(DATA_5M, "raw")
QFQ_5M = os.path.join(DATA_5M, "qfq")
RAW_15M = os.path.join(DATA_15M, "raw")
QFQ_15M = os.path.join(DATA_15M, "qfq")
PROGRESS_FILE = os.path.join(DATA_15M, "_progress.json")
RESULT_FILE = os.path.join(DATA_15M, "_result.json")
CURSOR_FILE = os.path.join(DATA_15M, "_cursor.json") # crash 前正在处理的股票(poison-pill 探测)
BLOCKLIST_FILE = os.path.join(DATA_15M, "_blocklist.json") # 累积 poison-pill 名单(跨重启)
START_TIME = "20250717"
END_TIME = dt.datetime.now().strftime("%Y%m%d")
# 实证预期:5m≈11670 bars、15m≈3890 bars(全板块一致,2025-07-17~today ~240 交易日)
EXPECTED_BARS_5M = 11670
EXPECTED_BARS_15M = 3890
MIN_BARS_RESUME_5M = 11400 # 断点续传阈值(允许 ±2.5%
MIN_BARS_RESUME_15M = 3888 # 15m 断点续传阈值(接近完整 3890,旧文件 600000=2040 会重下)
EXPECTED_BARS_PER_DAY_5M = 48 # A 股 4h 交易日 × 12 bar/h
EXPECTED_BARS_PER_DAY_15M = 16 # A 股 4h 交易日 × 4 bar/h
DOWNLOAD_BATCH = 200 # download_history_data2 批量大小(与 daily_update_xtdata 一致)
SLEEP_BETWEEN_BATCH = 2.0 # 批次间 sleep(不猛打券商后端)
SLEEP_BETWEEN_PHASES = 5.0 # 5m → 15m 轮次切换间 sleep
SLEEP_BETWEEN_STOCKS_15M = 0.2 # 15m 逐只 download 间隔(paced,绝不并发)
PROGRESS_FLUSH_EVERY = 200 # 进度文件刷新间隔(写盘阶段每 N 只)
FAILED_LIST_CAP = 200 # _progress.json 内 failed 列表截断
DISK_ALERT_GB = 20.0 # 磁盘 free < 此值告警
from xtquant import xtdata as xd
T0 = time.time()
def log(m: str) -> None:
print(f"[15M {time.time()-T0:.0f}s] {m}", flush=True)
def disk_free_gb() -> float:
"""C:\\ 剩余空间 GB(shutil 跨平台)。失败返回 -1。"""
try:
import shutil
return shutil.disk_usage("C:\\").free / (1024 ** 3)
except Exception: # noqa: BLE001
return -1.0
def count_parquet(directory: str, suffix: str) -> int:
if not os.path.isdir(directory):
return 0
try:
return len([f for f in os.listdir(directory) if f.endswith(f"_{suffix}.parquet")])
except Exception: # noqa: BLE001
return 0
def load_blocklist() -> set[str]:
"""读累积 poison-pill 名单。crash 后 wrapper 重启时,cursor 里那只会被加入此处。"""
if not os.path.exists(BLOCKLIST_FILE):
return set()
try:
with open(BLOCKLIST_FILE, "r", encoding="utf-8") as f:
return set(json.load(f).get("blocklist", []))
except Exception: # noqa: BLE001
return set()
def add_to_blocklist(code: str, reason: str) -> None:
"""把 poison-pill 股加入持久 blocklist(跨重启累积)。"""
bl = load_blocklist()
if code in bl:
return
bl.add(code)
payload = {
"blocklist": sorted(bl),
"last_added": code,
"last_reason": reason,
"last_added_at": dt.datetime.now().isoformat(),
"count": len(bl),
}
tmp = BLOCKLIST_FILE + ".tmp"
with open(tmp, "w", encoding="utf-8") as f:
json.dump(payload, f, ensure_ascii=False, indent=2)
os.replace(tmp, BLOCKLIST_FILE)
def write_cursor(code: str) -> None:
"""download_history_data 调用前写 cursor。crash 后重启时读 cursor → 该 code 是 poison-pill。"""
try:
tmp = CURSOR_FILE + ".tmp"
with open(tmp, "w", encoding="utf-8") as f:
json.dump({"processing": code, "ts": dt.datetime.now().isoformat()}, f,
ensure_ascii=False)
os.replace(tmp, CURSOR_FILE)
except Exception: # noqa: BLE001
pass
def clear_cursor() -> None:
"""单只处理完清除 cursor(成功路径)。"""
try:
if os.path.exists(CURSOR_FILE):
os.remove(CURSOR_FILE)
except Exception: # noqa: BLE001
pass
def read_stale_cursor() -> str | None:
"""启动时读 cursor。若存在 → 上次 crash 时正在处理的那只 = poison-pill。"""
if not os.path.exists(CURSOR_FILE):
return None
try:
with open(CURSOR_FILE, "r", encoding="utf-8") as f:
return json.load(f).get("processing")
except Exception: # noqa: BLE001
return None
def parse_dt_index(idx) -> pd.DatetimeIndex:
"""xtdata 时间索引(14 位 str/int 'yyyymmddHHMMSS' 或带毫秒)→ DatetimeIndex。"""
s = pd.Series([str(i) for i in idx]).str.slice(0, 14)
return pd.to_datetime(s, format="%Y%m%d%H%M%S", errors="coerce")
def to_pdf(df: pd.DataFrame) -> pd.DataFrame:
"""xtdata DataFrame → 统一 schema。volume ×100 手→股。"""
return pd.DataFrame({
"datetime": parse_dt_index(df.index),
"open": df["open"].astype(float).values,
"high": df["high"].astype(float).values,
"low": df["low"].astype(float).values,
"close": df["close"].astype(float).values,
"volume": (df["volume"].astype(float) * 100.0).values, # xtdata 手→股(memory 铁证)
}).dropna(subset=["datetime"]).sort_values("datetime").reset_index(drop=True)
def atomic_write(path: str, df: pd.DataFrame) -> None:
os.makedirs(os.path.dirname(path), exist_ok=True)
tmp = path + ".tmp"
df.to_parquet(tmp, index=False)
os.replace(tmp, path)
def write_progress(payload: dict[str, Any]) -> None:
payload = {**payload, "elapsed_sec": round(time.time() - T0, 1),
"updated": dt.datetime.now().isoformat()}
os.makedirs(os.path.dirname(PROGRESS_FILE), exist_ok=True)
tmp = PROGRESS_FILE + ".tmp"
with open(tmp, "w", encoding="utf-8") as f:
json.dump(payload, f, ensure_ascii=False, indent=2)
os.replace(tmp, PROGRESS_FILE)
def fetch(code: str, period: str, dividend_type: str) -> pd.DataFrame | None:
"""从 xtdata 本地缓存 get 出 bars。返回原始 xtdata DataFrame 或 None。"""
r = xd.get_market_data_ex([], [code], period=period, start_time=START_TIME,
end_time=END_TIME, dividend_type=dividend_type)
return r.get(code) if r else None
def batch_download(period: str, universe: list[str]) -> tuple[int, int]:
"""批量 download_history_data2 全市场 → xtdata 本地缓存。
成败不看 download 返回/回调,以下一步 get bars 为准(Main Agent 实证)。
返回 (n_batches, n_err_batches)。
"""
n_batches = (len(universe) + DOWNLOAD_BATCH - 1) // DOWNLOAD_BATCH
err_batches = 0
def _cb(data: Any, prog: float) -> None:
# callback 只 flush 进度,不做成败判定
if prog >= 100.0:
log(f" dl_{period} batch done (prog={prog:.1f})")
for bi in range(n_batches):
chunk = universe[bi * DOWNLOAD_BATCH:(bi + 1) * DOWNLOAD_BATCH]
try:
# download_history_data2 批量版,ret 可能为 None(正常),不看
xd.download_history_data2(chunk, period, START_TIME, END_TIME, _cb)
except Exception as e: # noqa: BLE001
err_batches += 1
log(f" dl_{period} batch#{bi} err: {e}(继续,单批 fail 不致命)")
if (bi + 1) % 5 == 0 or (bi + 1) == n_batches:
log(f" dl_{period} batch {bi+1}/{n_batches} ({(bi+1)*100/n_batches:.1f}%) err_batches={err_batches}")
time.sleep(SLEEP_BETWEEN_BATCH)
return (n_batches, err_batches)
def write_one_period(code: str, period: str, raw_dir: str, qfq_dir: str,
min_bars: int) -> dict[str, Any]:
"""download 阶段已完成,从 xtdata 缓存 get → 写 raw/qfq parquet。返回 dict。"""
raw_path = os.path.join(raw_dir, f"{code}_{period}.parquet")
qfq_path = os.path.join(qfq_dir, f"{code}_{period}.parquet")
res: dict[str, Any] = {"code": code, "period": period,
"raw_bars": 0, "qfq_bars": 0, "err": None, "skipped": False}
# 断点续传:两源都已存在且 bar 数 ≥ 阈值 → 跳过
if os.path.exists(raw_path) and os.path.exists(qfq_path):
try:
r_old = pd.read_parquet(raw_path)
q_old = pd.read_parquet(qfq_path)
if len(r_old) >= min_bars and len(q_old) >= min_bars:
res["raw_bars"] = len(r_old)
res["qfq_bars"] = len(q_old)
res["skipped"] = True
return res
except Exception:
pass # 文件损坏,下面重读重写
# get raw (none)
try:
rdf = fetch(code, period, "none")
if rdf is not None and len(rdf):
pdf = to_pdf(rdf)
if len(pdf):
atomic_write(raw_path, pdf)
res["raw_bars"] = len(pdf)
except Exception as e: # noqa: BLE001
res["err"] = f"raw:{e}"
# get qfq (front) —— 同份缓存读,不为 qfq 单独 download
try:
qdf = fetch(code, period, "front")
if qdf is not None and len(qdf):
pdf = to_pdf(qdf)
if len(pdf):
atomic_write(qfq_path, pdf)
res["qfq_bars"] = len(pdf)
except Exception as e: # noqa: BLE001
res["err"] = (res["err"] or "") + f" qfq:{e}"
return res
def write_phase(period: str, universe: list[str], raw_dir: str, qfq_dir: str,
min_bars: int, failed: list[dict], done_counter: dict[str, int],
counter_key: str) -> int:
"""逐只 get → 写 raw/qfq。失败(raw_bars=0)记 failed 不中断。返回 done 数。"""
done = 0
last_stock = "(start)"
fail_phase = 0
for i, code in enumerate(universe):
last_stock = code
try:
res = write_one_period(code, period, raw_dir, qfq_dir, min_bars)
if not res["skipped"] and res["raw_bars"] == 0 and res["qfq_bars"] == 0:
failed.append({"code": code, "period": period, "err": res["err"] or "no_data"})
fail_phase += 1
except Exception as e: # noqa: BLE001
failed.append({"code": code, "period": period, "err": str(e)})
fail_phase += 1
done = i + 1
done_counter[counter_key] = done
if done % PROGRESS_FLUSH_EVERY == 0 or done == len(universe):
write_progress({
"phase": f"write_{period}",
"total": len(universe),
"done_5m": done_counter.get("5m", 0),
"done_15m": done_counter.get("15m", 0),
"failed_count": len(failed),
"failed": failed[:FAILED_LIST_CAP],
"failed_truncated": len(failed) > FAILED_LIST_CAP,
"last_stock": last_stock,
"pct": round(done * 100 / len(universe), 2),
})
log(f" write_{period} {done}/{len(universe)} ({done*100/len(universe):.1f}%) "
f"phase_fail={fail_phase} total_fail={len(failed)} last={last_stock}")
return done
def download_and_write_one_15m(code: str, min_bars: int) -> dict[str, Any]:
"""15m 逐只 download_history_data(code,'15m',start,end) + get + 写 raw/qfq。
替代原 batch_download('15m')+write_one_period 组合。成败看 get bars >0
download_history_data 返回 None 是正常,Main Agent 实证)。
"""
raw_path = os.path.join(RAW_15M, f"{code}_15m.parquet")
qfq_path = os.path.join(QFQ_15M, f"{code}_15m.parquet")
res: dict[str, Any] = {"code": code, "period": "15m",
"raw_bars": 0, "qfq_bars": 0, "err": None, "skipped": False}
# 断点续传:两源已存在且 bars ≥ 阈值 → 跳过 download/get
if os.path.exists(raw_path) and os.path.exists(qfq_path):
try:
r_old = pd.read_parquet(raw_path)
q_old = pd.read_parquet(qfq_path)
if len(r_old) >= min_bars and len(q_old) >= min_bars:
res["raw_bars"] = len(r_old)
res["qfq_bars"] = len(q_old)
res["skipped"] = True
return res
except Exception:
pass # 文件损坏 → 重下
# 单只 downloadret=None 正常,看后续 get bars
try:
xd.download_history_data(code, "15m", START_TIME, END_TIME)
except Exception as e: # noqa: BLE001
res["err"] = f"dl:{e}"
# 不在此 return——download 报错也可能只是"已存在",继续 get
# get raw (none) → 写
try:
rdf = fetch(code, "15m", "none")
if rdf is not None and len(rdf):
pdf = to_pdf(rdf)
if len(pdf):
atomic_write(raw_path, pdf)
res["raw_bars"] = len(pdf)
except Exception as e: # noqa: BLE001
res["err"] = (res["err"] or "") + f" raw:{e}"
# get qfq (front) → 写(同份缓存读,不为 qfq 单独 download
try:
qdf = fetch(code, "15m", "front")
if qdf is not None and len(qdf):
pdf = to_pdf(qdf)
if len(pdf):
atomic_write(qfq_path, pdf)
res["qfq_bars"] = len(pdf)
except Exception as e: # noqa: BLE001
res["err"] = (res["err"] or "") + f" qfq:{e}"
return res
def phase_15m_per_stock(universe: list[str], failed: list[dict],
done_counter: dict[str, int]) -> int:
"""15m 逐只 download+write 循环。替代原 Phase 2 (batch_download) + Phase 2.5 (write_phase)。
单线程 paced(每只 sleep 0.2s),绝不并发。成败看 get bars >0。
【poison-pill 防护】
- 启动时读 _cursor.json:若存在 → 上次 crash 时正在处理的那只 = 嫌疑 poison-pill → 加入 blocklist
- 每只 download 前写 cursor,成功后清除
- blocklist 内的 code 跳过(记 failederr=blocked_poison_pill
"""
# 启动:检查 stale cursor → poison-pill
blocklist = load_blocklist()
stale = read_stale_cursor()
if stale:
log(f" ⚠️ detected stale cursor from crash: {stale} → adding to blocklist")
add_to_blocklist(stale, f"crash during download_history_data at {dt.datetime.now().isoformat()}")
blocklist = load_blocklist()
clear_cursor()
if blocklist:
log(f" blocklist ({len(blocklist)}): {sorted(blocklist)[:10]}{'...' if len(blocklist)>10 else ''}")
done = 0
fail_phase = 0
last_stock = "(start)"
for i, code in enumerate(universe):
last_stock = code
if code in blocklist:
failed.append({"code": code, "period": "15m", "err": "blocked_poison_pill"})
fail_phase += 1
done = i + 1
done_counter["15m"] = done
time.sleep(SLEEP_BETWEEN_STOCKS_15M)
continue
write_cursor(code) # 标记:即将处理此 code(crash 后重启可识别)
try:
res = download_and_write_one_15m(code, MIN_BARS_RESUME_15M)
if not res["skipped"] and res["raw_bars"] == 0 and res["qfq_bars"] == 0:
failed.append({"code": code, "period": "15m", "err": res["err"] or "no_data"})
fail_phase += 1
except Exception as e: # noqa: BLE001
failed.append({"code": code, "period": "15m", "err": str(e)})
fail_phase += 1
clear_cursor() # 成功路径:清除 cursor
done = i + 1
done_counter["15m"] = done
if done % PROGRESS_FLUSH_EVERY == 0 or done == len(universe):
dfree = disk_free_gb()
write_progress({
"phase": "write_15m",
"total": len(universe),
"done_5m": done_counter.get("5m", 0),
"done_15m": done,
"failed_count": len(failed),
"failed": failed[:FAILED_LIST_CAP],
"failed_truncated": len(failed) > FAILED_LIST_CAP,
"last_stock": last_stock,
"pct": round(done * 100 / len(universe), 2),
"disk_free_gb": round(dfree, 2),
"blocklist_size": len(blocklist),
})
log(f" write_15m {done}/{len(universe)} ({done*100/len(universe):.1f}%) "
f"phase_fail={fail_phase} total_fail={len(failed)} last={last_stock} "
f"disk_free={dfree:.1f}GB bl={len(blocklist)}")
if dfree >= 0 and dfree < DISK_ALERT_GB:
log(f" ⚠️ DISK LOW: {dfree:.1f}GB < {DISK_ALERT_GB}GB")
time.sleep(SLEEP_BETWEEN_STOCKS_15M)
return done
def validate(universe: list[str]) -> dict[str, Any]:
"""校验 5m + 15m 双周期 → 返回 result dict。"""
out: dict[str, Any] = {"periods": {}}
for period, raw_dir, qfq_dir, expected_bars_per_day, min_bars_floor in [
("5m", RAW_5M, QFQ_5M, EXPECTED_BARS_PER_DAY_5M, 5000),
("15m", RAW_15M, QFQ_15M, EXPECTED_BARS_PER_DAY_15M, 5000),
]:
raw_files = [f for f in os.listdir(raw_dir) if f.endswith(f"_{period}.parquet")] \
if os.path.isdir(raw_dir) else []
qfq_files = [f for f in os.listdir(qfq_dir) if f.endswith(f"_{period}.parquet")] \
if os.path.isdir(qfq_dir) else []
n_raw = len(raw_files)
n_qfq = len(qfq_files)
total_bars = 0
min_dt = None
max_dt = None
for f in raw_files:
try:
d = pd.read_parquet(os.path.join(raw_dir, f), columns=["datetime"])
if len(d):
total_bars += len(d)
lo = pd.to_datetime(d["datetime"]).min()
hi = pd.to_datetime(d["datetime"]).max()
if min_dt is None or lo < min_dt:
min_dt = lo
if max_dt is None or hi > max_dt:
max_dt = hi
except Exception:
continue
# 抽样 5 只校验完整性
samples = random.sample(raw_files, min(5, len(raw_files))) if raw_files else []
sample_results = []
complete_count = 0
for sf in samples:
try:
d = pd.read_parquet(os.path.join(raw_dir, sf))
d["_d"] = pd.to_datetime(d["datetime"]).dt.date
grp = d.groupby("_d").size()
incomplete_days = int((grp < expected_bars_per_day).sum())
sample_results.append({
"file": sf, "bars": int(len(d)), "days": int(len(grp)),
"incomplete_days": incomplete_days,
"first": str(d["datetime"].min()),
"last": str(d["datetime"].max()),
})
if incomplete_days == 0:
complete_count += 1
except Exception as e:
sample_results.append({"file": sf, "err": str(e)})
failed_stocks = [c for c in universe if f"{c}_{period}.parquet" not in set(raw_files)]
min_str = str(min_dt)[:19] if min_dt is not None else None
max_str = str(max_dt)[:19] if max_dt is not None else None
date_ok = False
if min_str and max_str:
earliest_limit = "2025-07-20" # 容忍首日 7/17~7/20 起步
latest_floor = (dt.datetime.now() - dt.timedelta(days=5)).strftime("%Y-%m-%d")
date_ok = min_str[:10] <= earliest_limit and max_str[:10] >= latest_floor
reasons = []
if n_raw < min_bars_floor:
reasons.append(f"raw 文件数 {n_raw} < {min_bars_floor}")
if len(samples) > 0 and complete_count < len(samples):
reasons.append(f"抽样完整 {complete_count}/{len(samples)}(有缺失日)")
if not date_ok:
reasons.append(f"日期范围异常: {min_str}~{max_str}")
verdict = "PASS" if (n_raw >= min_bars_floor
and (len(samples) == 0 or complete_count == len(samples))
and date_ok) else "FAIL"
out["periods"][period] = {
"verdict": verdict,
"fail_reasons": reasons,
"raw_files": n_raw,
"qfq_files": n_qfq,
"total_bars": int(total_bars),
"min_datetime": min_str,
"max_datetime": max_str,
"sample_size": len(samples),
"sample_complete_count": complete_count,
"sample": sample_results,
"failed_count": len(failed_stocks),
"failed_sample": failed_stocks[:30],
}
overall = "PASS" if all(p["verdict"] == "PASS" for p in out["periods"].values()) else "FAIL"
out["verdict"] = overall
out["universe_size"] = len(universe)
out["generated_at"] = dt.datetime.now().isoformat()
return out
def main() -> None:
log(f"START window={START_TIME}~{END_TIME} (先5m→后15m 两轮批量 download_history_data2)")
for d in (RAW_5M, QFQ_5M, RAW_15M, QFQ_15M):
os.makedirs(d, exist_ok=True)
# universe
try:
u = xd.get_stock_list_in_sector("沪深A股") or []
except Exception as e: # noqa: BLE001
log(f"FATAL: get_stock_list err: {e}")
os._exit(2)
if not u:
log("FATAL: empty universeminiQMT 未连?)")
os._exit(2)
log(f"universe={len(u)}(沪深A股,ETF 本轮未下)")
failed: list[dict] = []
done_counter = {"5m": 0, "15m": 0}
write_progress({
"phase": "init",
"total": len(u),
"done_5m": 0, "done_15m": 0,
"failed_count": 0, "failed": [], "last_stock": "(init)", "pct": 0.0,
})
# ===================== Phase 1 + 1.5: 5m(已下完则整体跳过) =====================
n5r = count_parquet(RAW_5M, "5m")
n5q = count_parquet(QFQ_5M, "5m")
skip_5m = (n5r >= len(u) and n5q >= len(u))
log(f"5m 现状: raw={n5r} qfq={n5q} universe={len(u)}"
f"{'SKIP(已下完)' if skip_5m else 'GO(需下)'} disk_free={disk_free_gb():.1f}GB")
if skip_5m:
log("PHASE 1+1.5 SKIP: 5m 已全量下完(断点续传保留)")
done_counter["5m"] = len(u)
else:
log("PHASE 1: 批量 download_history_data2 5m(成败看后续 get bars,不看 download ret")
nb1, errb1 = batch_download("5m", u)
log(f"PHASE 1 done: 5m batches={nb1} err_batches={errb1}")
log("PHASE 1.5: 逐只 get_market_data_ex 读 5m → 写 raw/qfq parquet")
write_phase("5m", u, RAW_5M, QFQ_5M, MIN_BARS_RESUME_5M, failed, done_counter, "5m")
log(f"PHASE 1.5 done: 5m write done={done_counter['5m']} total_fail={len(failed)}")
time.sleep(SLEEP_BETWEEN_PHASES)
# ===================== Phase 2: 15m 逐只 download+write(替代卡死的批量 download_history_data2 =====================
# Main Agent 实证:全市场批量 download_history_data2('15m', 5201 只一次)卡死 30min 0 产出;
# 但逐只 download_history_data(code,'15m',start,end) 实测稳定(600051/300001/688981/000001 均 3890 bars)。
log(f"PHASE 2: 逐只 download_history_data('15m') + get + write raw/qfq "
f"total={len(u)} disk_free={disk_free_gb():.1f}GB")
phase_15m_per_stock(u, failed, done_counter)
log(f"PHASE 2 done: 15m write done={done_counter['15m']} total_fail={len(failed)} "
f"disk_free={disk_free_gb():.1f}GB")
write_progress({
"phase": "validate_pending",
"total": len(u),
"done_5m": done_counter["5m"], "done_15m": done_counter["15m"],
"failed_count": len(failed), "failed": failed[:FAILED_LIST_CAP],
"failed_truncated": len(failed) > FAILED_LIST_CAP,
"last_stock": "(validate)", "pct": 100.0,
})
# ===================== Phase 3: 校验 =====================
log("PHASE 3: 校验 5m + 15m 双周期")
result = validate(u)
tmp = RESULT_FILE + ".tmp"
with open(tmp, "w", encoding="utf-8") as f:
json.dump(result, f, ensure_ascii=False, indent=2)
os.replace(tmp, RESULT_FILE)
write_progress({
"phase": f"done:{result['verdict']}",
"total": len(u),
"done_5m": done_counter["5m"], "done_15m": done_counter["15m"],
"failed_count": len(failed), "failed": failed[:FAILED_LIST_CAP],
"failed_truncated": len(failed) > FAILED_LIST_CAP,
"last_stock": "(done)", "pct": 100.0,
})
for period, info in result["periods"].items():
log(f"VERDICT[{period}]={info['verdict']} raw={info['raw_files']} qfq={info['qfq_files']} "
f"bars={info['total_bars']} sample_complete={info['sample_complete_count']}/{info['sample_size']} "
f"failed={info['failed_count']}")
if info["min_datetime"] and info["max_datetime"]:
log(f" range[{period}] {info['min_datetime']} ~ {info['max_datetime']}")
if info["fail_reasons"]:
log(f" fail_reasons[{period}]: {info['fail_reasons']}")
log(f"OVERALL VERDICT={result['verdict']} total_failed={len(failed)}")
sys.stdout.flush()
os._exit(0 if result["verdict"] == "PASS" else 1)
if __name__ == "__main__":
main()
@@ -0,0 +1,240 @@
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""高效导入 5m / 15m 分钟 Parquet 到 vnpy SQLite dbbardata。
仿 import_vnpy_daily_fast.py 结构parse_filename + pandas 向量化 + batch
INSERT OR REPLACE + 末尾 dbbaroverview 刷新改成读 minute_5/minute_15 raw
parquet 灌到 dbbardata
口径
- 只灌 rawqfq 不灌同日线机制复权在读时按 dividend_type
- 文件名:<code>.<SH|SZ>_5m.parquet / _15m.parquet实证 (sh|sz)<code>
- interval:'5m' 文件灌 '5m''15m' 文件灌 '15m'字符串字段
- volume 已是股parquet ×100 别再 ×
- turnover 若无则 0.0open_interest=0.0
- datetime:parquet datetime64[ns] astype(str) 'YYYY-MM-DD HH:MM:SS'同日线机制
幂等INSERT OR REPLACE(symbol,exchange,datetime,interval unique重跑安全
进度 chunk flush print断点续传可选全量灌一次
用法VPS
set VNPY_DB_PATH=C:\\sanguo_vnpy_v2\\data\\quant_trading.db
set MINUTE_5_DIR=C:\\sanguo_vnpy_v2\\data\\minute_5\\raw
set MINUTE_15_DIR=C:\\sanguo_vnpy_v2\\data\\minute_15\\raw
C:\\Python310\\python.exe -X utf8 import_vnpy_minute_fast.py
可选参数 --only 5m / --only 15m 只灌一个周期
退出码 0=完成1=致命错误
"""
import os
import re
import sys
import time
import sqlite3
from pathlib import Path
import pandas as pd
DB_PATH = os.environ.get(
'VNPY_DB_PATH', r'C:\sanguo_vnpy_v2\data\quant_trading.db'
)
MINUTE_5_DIR = os.environ.get(
'MINUTE_5_DIR', r'C:\sanguo_vnpy_v2\data\minute_5\raw'
)
MINUTE_15_DIR = os.environ.get(
'MINUTE_15_DIR', r'C:\sanguo_vnpy_v2\data\minute_15\raw'
)
# chunk 内文件数(向量化合并 + 批量 INSERT 的平衡点;5m 每 chunk ~120w 行)
# 经验:500 files/chunk 在 SSH-detached 环境下疑似资源约束 OOM-kill;降到 100 稳妥。
CHUNK_FILES = int(os.environ.get('MINUTE_CHUNK_FILES', '100'))
BATCH_INSERT = 20000
# 文件名 <code>.<EXC>_5m.parquet(实证格式)
_FNAME_RE = re.compile(r'^(\d{6})\.(SH|SZ)_\d+m\.parquet$', re.IGNORECASE)
def parse_filename(filename: str):
"""<code>.<SH|SZ>_Nm.parquet → (code, exchange_name) or (None, None)."""
m = _FNAME_RE.match(filename)
if not m:
return None, None
code, exc = m.groups()
exchange = 'SSE' if exc.upper() == 'SH' else 'SZSE'
return code, exchange
def list_parquets(raw_dir: str, period_suffix: str):
"""raw_dir 下匹配 <code>.<exc>_<period>.parquet 的文件。"""
p = Path(raw_dir)
if not p.exists():
return []
return sorted(p.glob(f'*_{period_suffix}.parquet'))
def import_one_chunk(conn, files, interval: str, chunk_idx: int, total_chunks: int):
"""向量化读 → 合并 → batch INSERT OR REPLACE。返回 (n_files_ok, n_rows).
内存优化 N 个文件 sub-chunk避免单 chunk concat/tolist 占用过大
"""
c = conn.cursor()
n_ok_total = 0
n_rows_total = 0
SUB_BATCH = 25 # 每 25 文件 concat→insert→释放,控制峰值
for si in range(0, len(files), SUB_BATCH):
sub = files[si:si + SUB_BATCH]
dfs = []
n_ok = 0
for f in sub:
code, exchange = parse_filename(f.name)
if code is None:
continue
try:
df = pd.read_parquet(
f, columns=['datetime', 'open', 'high', 'low', 'close', 'volume']
)
except Exception:
continue
if df.empty:
continue
df['symbol'] = code
df['exchange'] = exchange
dfs.append(df)
n_ok += 1
if not dfs:
continue
combined = pd.concat(dfs, ignore_index=True)
del dfs
combined['datetime'] = combined['datetime'].astype(str)
combined['interval'] = interval
combined['open_interest'] = 0.0
combined = combined.rename(columns={
'open': 'open_price', 'high': 'high_price',
'low': 'low_price', 'close': 'close_price',
})
if 'turnover' not in combined.columns:
combined['turnover'] = 0.0
for col in ('volume', 'turnover', 'open_price',
'high_price', 'low_price', 'close_price'):
combined[col] = combined[col].fillna(0.0).astype(float)
values = combined[[
'symbol', 'exchange', 'datetime', 'interval', 'volume', 'turnover',
'open_interest', 'open_price', 'high_price', 'low_price', 'close_price',
]].values.tolist()
del combined
for i in range(0, len(values), BATCH_INSERT):
c.executemany(
'''INSERT OR REPLACE INTO dbbardata
(symbol,exchange,datetime,interval,volume,turnover,open_interest,
open_price,high_price,low_price,close_price)
VALUES (?,?,?,?,?,?,?,?,?,?,?)''',
values[i:i + BATCH_INSERT],
)
conn.commit()
n_ok_total += n_ok
n_rows_total += len(values)
del values
return n_ok_total, n_rows_total
def import_interval(conn, interval: str, raw_dir: str, period_suffix: str):
files = list_parquets(raw_dir, period_suffix)
n_total = len(files)
if n_total == 0:
print(f'[{interval}] no parquet under {raw_dir}', flush=True)
return 0, 0
n_chunks = (n_total + CHUNK_FILES - 1) // CHUNK_FILES
print(f'[{interval}] {n_total} files, {n_chunks} chunks, dir={raw_dir}',
flush=True)
total_rows = 0
total_files = 0
t0 = time.time()
for ci in range(n_chunks):
chunk_files = files[ci * CHUNK_FILES:(ci + 1) * CHUNK_FILES]
nf, nrows = import_one_chunk(conn, chunk_files, interval, ci, n_chunks)
total_files += nf
total_rows += nrows
elapsed = time.time() - t0
print(f'[{interval}] chunk {ci+1}/{n_chunks} '
f'files_done={total_files}/{n_total} '
f'rows={total_rows} elapsed={elapsed:.1f}s '
f'({total_rows/max(elapsed,1):.0f} rows/s)', flush=True)
return total_files, total_rows
def update_overview(conn):
"""刷新 dbbaroverview(按 symbol,exchange,interval 分组聚合,含 5m/15m)。"""
c = conn.cursor()
c.execute(
'''INSERT OR REPLACE INTO dbbaroverview
(symbol,exchange,interval,count,start,end)
SELECT symbol,exchange,interval,COUNT(*),MIN(datetime),MAX(datetime)
FROM dbbardata GROUP BY symbol,exchange,interval'''
)
conn.commit()
def report_db_counts(conn):
c = conn.cursor()
rows = c.execute(
'SELECT interval, COUNT(*) FROM dbbardata GROUP BY interval'
).fetchall()
for iv, n in rows:
print(f' interval={iv} count={n}', flush=True)
for iv in ('5m', '15m'):
row = c.execute(
'SELECT MAX(datetime), MIN(datetime), COUNT(*) '
'FROM dbbardata WHERE interval=?', (iv,)
).fetchone()
if row and row[0]:
print(f' [{iv}] range={row[1]} ~ {row[0]} rows={row[2]}', flush=True)
def main():
only = None
for i, arg in enumerate(sys.argv):
if arg == '--only' and i + 1 < len(sys.argv):
only = sys.argv[i + 1].lower()
print(f'Import minute parquet → DB: {DB_PATH}', flush=True)
print(f' MINUTE_5_DIR={MINUTE_5_DIR}', flush=True)
print(f' MINUTE_15_DIR={MINUTE_15_DIR}', flush=True)
print(f' only={only}', flush=True)
if not os.path.exists(DB_PATH):
print(f'FATAL: DB not found: {DB_PATH}', flush=True)
sys.exit(1)
conn = sqlite3.connect(DB_PATH)
t_start = time.time()
grand_rows = 0
if only in (None, '5m'):
nf, nr = import_interval(conn, '5m', MINUTE_5_DIR, '5m')
grand_rows += nr
print(f'[5m] DONE files={nf} rows={nr}', flush=True)
if only in (None, '15m'):
nf, nr = import_interval(conn, '15m', MINUTE_15_DIR, '15m')
grand_rows += nr
print(f'[15m] DONE files={nf} rows={nr}', flush=True)
print(f'Updating dbbaroverview ...', flush=True)
t_ov = time.time()
update_overview(conn)
print(f' overview updated ({time.time()-t_ov:.1f}s)', flush=True)
elapsed = time.time() - t_start
print(f'\nDone in {elapsed:.1f}s ({elapsed/60:.1f}min) total_rows_added={grand_rows}',
flush=True)
print('Final DB counts:', flush=True)
report_db_counts(conn)
conn.close()
sys.exit(0)
if __name__ == '__main__':
main()
@@ -0,0 +1,40 @@
$ErrorActionPreference = 'Continue'
$ProgressPreference = 'SilentlyContinue'
$maxRetries = 15
$script = 'C:\sanguo_vnpy_v2\scripts\data_platform\download_15m_xtdata.py'
$resultFile = 'C:\sanguo_vnpy_v2\data\minute_15\_result.json'
$logFile = 'C:\sanguo_vnpy_v2\data\minute_15_download.log'
$errFile = 'C:\sanguo_vnpy_v2\data\minute_15_download.err.log'
$wrapperLog = 'C:\sanguo_vnpy_v2\data\minute_15_wrapper.log'
function Log-W($m) {
$line = "$(Get-Date -Format 'yyyy-MM-dd HH:mm:ss') $m"
Add-Content -Path $wrapperLog -Value $line -Encoding UTF8
}
Log-W "WRAPPER START maxRetries=$maxRetries"
for ($i = 1; $i -le $maxRetries; $i++) {
Log-W "ATTEMPT $i/$maxRetries starting python"
if (Test-Path $errFile) { Clear-Content $errFile -EA SilentlyContinue }
try {
$p = Start-Process -FilePath 'C:\Python310\python.exe' `
-ArgumentList '-X utf8', $script `
-WindowStyle Hidden `
-RedirectStandardOutput $logFile `
-RedirectStandardError $errFile `
-Wait -PassThru
$ec = if ($null -ne $p.ExitCode) { $p.ExitCode } else { -999 }
} catch {
$ec = -998
Log-W "Start-Process exception: $_"
}
Log-W "ATTEMPT $i exited code=$ec"
if (Test-Path $resultFile) {
Log-W "_result.json EXISTS - DONE"
break
}
Log-W "no result yet, sleeping 8s before retry..."
Start-Sleep -Seconds 8
}
Log-W "WRAPPER EXIT"
+262
View File
@@ -0,0 +1,262 @@
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""灌后校验:parquet 日线 ↔ quant_trading.db dbbardata 三项一致性检查。
独立可跑读两个环境变量
VNPY_DB_PATH 默认 C:\\sanguo_vnpy_v2\\data\\quant_trading.db
DAILY_DIR 默认 C:\\sanguo_vnpy_v2\\data\\raw
三项校验
A 最新日期相等: max(date) over raw parquet == max(datetime) where interval='d' from DB
B 抽查 3 只股 OHLCV 逐值相等: 600519(SSE)/000001(SZSE)/000858(SZSE),
parquet 最新 5 vs DB datetime 倒序 limit 5, date 对齐逐值比对 (abs<1e-4)
C 当天数据合理性: DB date==parquet_max interval='d' 的所有行,
close_price<=0 high_price<low_priceclose_price null
退出码: PASS=0 / FAIL=1 stdlib + pandas (VPS 已装)
Windows GBK 终端需 `python -X utf8 validate_import.py` 调用
"""
import os
import sys
import sqlite3
from pathlib import Path
import pandas as pd
DB_PATH = os.environ.get('VNPY_DB_PATH', r'C:\sanguo_vnpy_v2\data\quant_trading.db')
DAILY_DIR = os.environ.get('DAILY_DIR', r'C:\sanguo_vnpy_v2\data\raw')
# (symbol_code, exchange, filename_prefix) 三只蓝筹股做抽查
SAMPLES = [
('600519', 'SSE', 'sh'),
('000001', 'SZSE', 'sz'),
('000858', 'SZSE', 'sz'),
]
def norm_date_str(s):
"""统一成 'YYYY-MM-DD' 字符串。接受 Timestamp / datetime / str。"""
ts = pd.Timestamp(s)
return ts.strftime('%Y-%m-%d')
def find_parquet_max_date():
"""扫最近年份目录下所有 raw parquet,取最大 date。
只扫最新有数据的年份目录~5000 文件 × 仅读 date 10s
返回 (max_date_str or None, n_files_scanned)"""
max_date = None
n_files = 0
year_dirs = sorted(
[d for d in Path(DAILY_DIR).glob('*') if d.is_dir() and d.name.isdigit()],
key=lambda d: d.name,
reverse=True,
)
for year_dir in year_dirs[:2]: # 最多看最近 2 个年份目录
local_max = None
for f in year_dir.glob('*.parquet'):
n_files += 1
try:
df = pd.read_parquet(f, columns=['date'])
if df.empty:
continue
cur = df['date'].max()
if local_max is None or cur > local_max:
local_max = cur
except Exception:
continue
if local_max is not None:
max_date = local_max
break # 当前年份有数据就不再看上一年
return (norm_date_str(max_date) if max_date is not None else None, n_files)
def check_a(parquet_max, db_max):
"""A: 最新日期相等"""
p = parquet_max
d = db_max[:10] if db_max else None
ok = (p is not None) and (d is not None) and (p == d)
return ok, {'parquet_max': p, 'db_max': d}
def check_b(conn):
"""B: 3 只股 OHLCV 逐值比对"""
detail = {}
all_ok = True
for code, exchange, prefix in SAMPLES:
# 从最新年份目录往前找 parquet
year_dirs = sorted(
[d for d in Path(DAILY_DIR).glob('*') if d.is_dir() and d.name.isdigit()],
key=lambda d: d.name,
reverse=True,
)
parquet_path = None
for yd in year_dirs:
p = yd / f'{prefix}{code}_daily.parquet'
if p.exists():
parquet_path = p
break
if parquet_path is None:
detail[code] = {'status': 'MISSING_PARQUET'}
all_ok = False
continue
try:
pdf = pd.read_parquet(parquet_path).sort_values('date').tail(5).reset_index(drop=True)
except Exception as e:
detail[code] = {'status': f'PARQUET_READ_ERR: {e}'}
all_ok = False
continue
try:
rows = conn.execute(
"SELECT datetime, open_price, high_price, low_price, close_price, volume "
"FROM dbbardata WHERE symbol=? AND exchange=? AND interval='d' "
"ORDER BY datetime DESC LIMIT 5",
(code, exchange),
).fetchall()
except Exception as e:
detail[code] = {'status': f'DB_QUERY_ERR: {e}'}
all_ok = False
continue
if not rows:
detail[code] = {'status': 'DB_EMPTY'}
all_ok = False
continue
db_pdf = pd.DataFrame(rows, columns=['datetime', 'open', 'high', 'low', 'close', 'volume'])
db_pdf = db_pdf.sort_values('datetime').reset_index(drop=True)
db_pdf['date'] = db_pdf['datetime'].str[:10]
pdf = pdf.copy()
pdf['date'] = pdf['date'].astype(str).str[:10]
# 取共同尾部 n 行
n = min(len(pdf), len(db_pdf))
if n == 0:
detail[code] = {'status': 'EMPTY_AFTER_ALIGN'}
all_ok = False
continue
p_tail = pdf.tail(n).reset_index(drop=True)
d_tail = db_pdf.tail(n).reset_index(drop=True)
dates_match = (p_tail['date'].values == d_tail['date'].values).all()
cols = ['open', 'high', 'low', 'close', 'volume']
vals_match = True
bad_col = None
for c in cols:
try:
pv = p_tail[c].astype(float).values
dv = d_tail[c].astype(float).values
max_diff = float((abs(pv - dv)).max())
if max_diff > 1e-4:
vals_match = False
bad_col = f'{c}(max_diff={max_diff})'
break
except Exception as e:
vals_match = False
bad_col = f'{c}: {e}'
break
ok = bool(dates_match and vals_match)
if not ok:
all_ok = False
detail[code] = {
'status': 'PASS' if ok else 'FAIL',
'n_rows': n,
'dates_match': bool(dates_match),
'vals_match': bool(vals_match),
'bad_col': bad_col,
'latest_date': p_tail['date'].iloc[-1] if len(p_tail) else None,
'latest_close': float(p_tail['close'].iloc[-1]) if len(p_tail) else None,
}
return all_ok, detail
def check_c(conn, parquet_max_str):
"""C: 当天数据合理性(DB 中 date==parquet_max 的所有 interval='d' 行)"""
if not parquet_max_str:
return False, {'reason': 'no parquet_max'}
pattern = parquet_max_str + '%'
try:
total = conn.execute(
"SELECT COUNT(*) FROM dbbardata WHERE interval='d' AND datetime LIKE ?",
(pattern,),
).fetchone()[0]
except Exception as e:
return False, {'reason': f'DB_QUERY_ERR: {e}'}
if total == 0:
return False, {'reason': f'no rows for {parquet_max_str}'}
try:
bad_close = conn.execute(
"SELECT COUNT(*) FROM dbbardata WHERE interval='d' AND datetime LIKE ? "
"AND (close_price IS NULL OR close_price <= 0)",
(pattern,),
).fetchone()[0]
bad_hl = conn.execute(
"SELECT COUNT(*) FROM dbbardata WHERE interval='d' AND datetime LIKE ? "
"AND high_price < low_price",
(pattern,),
).fetchone()[0]
except Exception as e:
return False, {'reason': f'DB_QUERY_ERR: {e}'}
ok = (bad_close == 0) and (bad_hl == 0)
return ok, {
'date': parquet_max_str,
'total_rows': total,
'bad_close': bad_close,
'bad_high_low': bad_hl,
}
def main():
print(f'[VALIDATE] DB={DB_PATH}')
print(f'[VALIDATE] DAILY_DIR={DAILY_DIR}')
if not os.path.exists(DB_PATH):
print(f'[VALIDATE][FAIL] DB not found: {DB_PATH}')
sys.exit(1)
if not os.path.exists(DAILY_DIR):
print(f'[VALIDATE][FAIL] DAILY_DIR not found: {DAILY_DIR}')
sys.exit(1)
parquet_max, n_files = find_parquet_max_date()
if parquet_max is None:
print(f'[VALIDATE][FAIL] no parquet under {DAILY_DIR} (scanned {n_files} files)')
sys.exit(1)
print(f'[VALIDATE] scanned {n_files} parquet files, parquet_max={parquet_max}')
conn = sqlite3.connect(DB_PATH)
try:
row = conn.execute("SELECT MAX(datetime) FROM dbbardata WHERE interval='d'").fetchone()
db_max = row[0] if row else None
except Exception as e:
print(f'[VALIDATE][FAIL] DB query max(datetime) err: {e}')
conn.close()
sys.exit(1)
print(f'[VALIDATE] db_max={db_max}')
# A
a_ok, a_detail = check_a(parquet_max, db_max)
print(f'[VALIDATE][A] {"PASS" if a_ok else "FAIL"} parquet_max={a_detail["parquet_max"]} db_max={a_detail["db_max"]}')
# B
b_ok, b_detail = check_b(conn)
print(f'[VALIDATE][B] {"PASS" if b_ok else "FAIL"}')
for code, info in b_detail.items():
if info.get('status') == 'PASS':
print(f' {code}: PASS n={info["n_rows"]} latest={info["latest_date"]} close={info["latest_close"]}')
else:
print(f' {code}: {info}')
# C
c_ok, c_detail = check_c(conn, parquet_max)
print(f'[VALIDATE][C] {"PASS" if c_ok else "FAIL"} {c_detail}')
conn.close()
overall = a_ok and b_ok and c_ok
print(f'[VALIDATE] {"PASS" if overall else "FAIL"} '
f'(A={"PASS" if a_ok else "FAIL"} B={"PASS" if b_ok else "FAIL"} C={"PASS" if c_ok else "FAIL"})')
sys.exit(0 if overall else 1)
if __name__ == '__main__':
main()