feat(data): build_daily_from_xtdata 全市场日线从xtdata重建到staging
miniQMT的xtdata作源(已证价格与parquet一致/T+0/深历史),重建VPS全市场日线。
- 只写 staging(DATA/_staging_xtdata),不动主库,验后再swap
- download一次raw,读两次(dividend_type none=raw/front=qfq,复权读时应用)
- volume×100(手→股),列 date/open/high/low/close/volume 对齐 datareader
- 单线程paced(分批200+sleep),per-stock写不全量进内存,避全市场崩
- 布局 {raw,qfq}/<year>/{sh|sz}{sym}_daily.parquet,沪深A股~5201只
This commit is contained in:
@@ -0,0 +1,108 @@
|
||||
#!/usr/bin/env python3
|
||||
# -*- coding: utf-8 -*-
|
||||
"""全市场日线:从 xtdata(miniQMT) 重建到 **staging**(不动主库)。
|
||||
|
||||
VPS 跑(需 miniQMT 常驻 + xtquant)。单线程 paced、per-stock 写、内存安全
|
||||
(不全量 load 进内存,逐只写盘——避 macOS Jetsman 那种全市场崩)。
|
||||
|
||||
输出:DATA/_staging_xtdata/{raw,qfq}/<year>/{sh|sz}{symbol}_daily.parquet
|
||||
列:date,open,high,low,close,volume(datareader::_row_to_bar 只读这 6 列)。
|
||||
volume ×100(xtdata 按"手",parquet 按"股")。
|
||||
下载策略:download_history_data2 下一次(raw 1d),get_market_data_ex 读两次
|
||||
(dividend_type=none→raw / front→qfq,复权是读时应用)。
|
||||
|
||||
跑完人工验 staging(count/日期/抽样),再单独 swap 进 data/{raw,qfq}。
|
||||
"""
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
import pandas as pd
|
||||
|
||||
ROOT = r"C:\sanguo_vnpy_v2"
|
||||
DATA = os.path.join(ROOT, "data")
|
||||
STAGING = os.path.join(DATA, "_staging_xtdata")
|
||||
START, END = "20100101", "20260715"
|
||||
|
||||
from xtquant import xtdata as xd
|
||||
|
||||
T0 = time.time()
|
||||
|
||||
|
||||
def log(m):
|
||||
print(f"[BUILD {time.time()-T0:.0f}s] {m}", flush=True)
|
||||
|
||||
|
||||
def prefix_of(sym):
|
||||
return "sh" if sym[:2] in ("60", "68", "51", "56", "58") else "sz"
|
||||
|
||||
|
||||
log("start")
|
||||
u = xd.get_stock_list_in_sector("沪深A股") or []
|
||||
log(f"universe={len(u)} sample={u[:3]}")
|
||||
|
||||
# 1. 下载一次(raw 日线),分批 paced(别猛打券商后端)
|
||||
BATCH = 200
|
||||
for i in range(0, len(u), BATCH):
|
||||
chunk = u[i:i + BATCH]
|
||||
try:
|
||||
xd.download_history_data2(chunk, "1d", START, END, lambda d, p: None)
|
||||
except Exception as e: # noqa: BLE001
|
||||
log(f"dl batch@{i} err: {e}")
|
||||
time.sleep(1.0)
|
||||
if i % 1000 == 0:
|
||||
log(f"download {i}/{len(u)}")
|
||||
log("download phase done")
|
||||
|
||||
|
||||
def fetch(code, dt):
|
||||
r = xd.get_market_data_ex([], [code], period="1d", start_time=START, end_time=END, dividend_type=dt)
|
||||
return r.get(code) if r else None
|
||||
|
||||
|
||||
def write_series(code, dt, kind):
|
||||
sym = code.split(".")[0]
|
||||
prefix = prefix_of(sym)
|
||||
df = fetch(code, dt)
|
||||
if df is None or not len(df):
|
||||
return 0
|
||||
pdf = pd.DataFrame({
|
||||
"date": pd.to_datetime([str(i)[:8] for i in 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).values, # 手→股
|
||||
})
|
||||
pdf["_y"] = pdf["date"].dt.year
|
||||
for yr, sub in pdf.groupby("_y"):
|
||||
d = os.path.join(STAGING, kind, str(int(yr)))
|
||||
os.makedirs(d, exist_ok=True)
|
||||
sub.drop(columns=["_y"]).to_parquet(
|
||||
os.path.join(d, f"{prefix}{sym}_daily.parquet"), index=False)
|
||||
return len(pdf)
|
||||
|
||||
|
||||
# 2. 读两次(raw/qfq)逐只写 staging
|
||||
ok = fail = 0
|
||||
nraw = nqfq = 0
|
||||
for i, code in enumerate(u):
|
||||
try:
|
||||
nr = write_series(code, "none", "raw")
|
||||
nq = write_series(code, "front", "qfq")
|
||||
nraw += nr
|
||||
nqfq += nq
|
||||
if nr or nq:
|
||||
ok += 1
|
||||
else:
|
||||
fail += 1
|
||||
except Exception as e: # noqa: BLE001
|
||||
fail += 1
|
||||
if fail <= 5:
|
||||
log(f"write err {code}: {e}")
|
||||
if (i + 1) % 500 == 0:
|
||||
log(f"write {i+1}/{len(u)} ok={ok} fail={fail}")
|
||||
|
||||
log(f"DONE stocks_ok={ok} fail={fail} raw_bars={nraw} qfq_bars={nqfq}")
|
||||
log(f"staging={STAGING}")
|
||||
sys.stdout.flush()
|
||||
os._exit(0)
|
||||
Reference in New Issue
Block a user