diff --git a/scripts/data_platform/_launch_15m.bat b/scripts/data_platform/_launch_15m.bat new file mode 100644 index 0000000..0a15cbb --- /dev/null +++ b/scripts/data_platform/_launch_15m.bat @@ -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% diff --git a/scripts/data_platform/_run_build.ps1 b/scripts/data_platform/_run_build.ps1 new file mode 100644 index 0000000..ff2f24b --- /dev/null +++ b/scripts/data_platform/_run_build.ps1 @@ -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 diff --git a/scripts/data_platform/_run_daily.ps1 b/scripts/data_platform/_run_daily.ps1 new file mode 100644 index 0000000..a5088a4 --- /dev/null +++ b/scripts/data_platform/_run_daily.ps1 @@ -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 增量下载脚本(inline,xtquant 单线程 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 + # 单只 download(ret=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 diff --git a/scripts/data_platform/_run_import.ps1 b/scripts/data_platform/_run_import.ps1 new file mode 100644 index 0000000..0f47052 --- /dev/null +++ b/scripts/data_platform/_run_import.ps1 @@ -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 diff --git a/scripts/data_platform/download_15m_xtdata.py b/scripts/data_platform/download_15m_xtdata.py new file mode 100644 index 0000000..c603858 --- /dev/null +++ b/scripts/data_platform/download_15m_xtdata.py @@ -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.py:download_history_data2 批量 → get_market_data_ex 分别 + 读 none/front 两份)。 + +【两轮下载顺序】 + Phase 1: 全市场 download_history_data2 批量下 5m(→ xtdata 本地缓存) + Phase 1.5: 逐只 get_market_data_ex 读 5m(raw + qfq)→ 写 parquet + Phase 2: 全市场 download_history_data2 批量下 15m(依赖 5m 基础) + Phase 2.5: 逐只 get_market_data_ex 读 15m(raw + qfq)→ 写 parquet + Phase 3: 校验 → 写 _result.json + +【存储】C:\\sanguo_vnpy_v2\\data\\minute_5\\{raw,qfq}\\_5m.parquet + C:\\sanguo_vnpy_v2\\data\\minute_15\\{raw,qfq}\\_15m.parquet + schema: datetime,open,high,low,close,volume(volume ×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 # 文件损坏 → 重下 + + # 单只 download(ret=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 跳过(记 failed,err=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 universe(miniQMT 未连?)") + 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() diff --git a/scripts/data_platform/import_vnpy_minute_fast.py b/scripts/data_platform/import_vnpy_minute_fast.py new file mode 100644 index 0000000..cffe5a2 --- /dev/null +++ b/scripts/data_platform/import_vnpy_minute_fast.py @@ -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。 + +【口径】 +- 只灌 raw(qfq 不灌——同日线机制,复权在读时按 dividend_type 转) +- 文件名:._5m.parquet / _15m.parquet(实证,非 (sh|sz)) +- interval:'5m' 文件灌 '5m'、'15m' 文件灌 '15m'(字符串字段) +- volume 已是股(parquet 里 ×100 过,别再 ×) +- turnover 若无则 0.0;open_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 + +# 文件名 ._5m.parquet(实证格式) +_FNAME_RE = re.compile(r'^(\d{6})\.(SH|SZ)_\d+m\.parquet$', re.IGNORECASE) + + +def parse_filename(filename: str): + """._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 下匹配 ._.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() diff --git a/scripts/data_platform/relaunch_15m_wrapper.ps1 b/scripts/data_platform/relaunch_15m_wrapper.ps1 new file mode 100644 index 0000000..ce10202 --- /dev/null +++ b/scripts/data_platform/relaunch_15m_wrapper.ps1 @@ -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" diff --git a/scripts/data_platform/validate_import.py b/scripts/data_platform/validate_import.py new file mode 100644 index 0000000..ac38630 --- /dev/null +++ b/scripts/data_platform/validate_import.py @@ -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 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()