ee27313644
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。
273 lines
12 KiB
PowerShell
273 lines
12 KiB
PowerShell
# 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
|