Files
sanguo_vnpy_v2/scripts/data_platform/daily_update_xtdata.py
T
claude_dev b966fce1e9 feat(data): 导入脚本 env 化 + xtdata 日线增量脚本
- import_vnpy_daily_fast.py: DB_PATH/DAILY_DIR 读 env(VNPY_DB_PATH/DAILY_DIR, Windows用正斜杠); amount 列缺失时容错填0
- daily_update_xtdata.py: 新增, 全市场日线从 miniQMT xtdata 本地缓存增量下载(零漂移), 部署文档 §4/§8 引用
2026-07-17 08:24:01 +08:00

170 lines
5.7 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""全市场日线每日增量补全(xtdata → 主库 data/{raw,qfq})。VPS 跑(需 miniQMT 常驻)。
设计:
- 增量窗口 LOOKBACK_DAYS 天(默认 30):覆盖周末/节假日 gap,不重写历史,省时省 IO。
- 先 download_history_data2 批量下载增量到 xtdata 本地缓存(单线程 paced,别猛打券商后端)。
- 再逐只 get_market_data_ex 读 raw(dividend_type=none)/qfq(front),与现有主库 parquet 按 date 去重
合并,原子写回对应年份文件(跨年自动落到两个年份目录)。
- 轻量验证:OHLC NaN 检查;失败只 warn 不阻断。
- 幂等:重复跑同一天不重复写(drop_duplicates keep='last')。
列与 build_daily_from_xtdata.py 完全一致:date,open,high,low,close,volumevolume×100 手→股)。
datareader._row_to_bar 只读这 6 列。
用法(VPS):
C:\\Python310\\python.exe -X utf8 daily_update_xtdata.py
可选环境变量:LOOKBACK_DAYS(默认 30)。
退出码:0=全部成功;1=有失败但流程完成;2=致命错误(universe 拉不到等)。
"""
import os
import sys
import time
import datetime as dt
import pandas as pd
ROOT = r"C:\sanguo_vnpy_v2"
DATA = os.path.join(ROOT, "data")
LOOKBACK = int(os.environ.get("LOOKBACK_DAYS", "30"))
from xtquant import xtdata as xd
T0 = time.time()
def log(m):
print(f"[DAILY {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"
def today_str():
return dt.datetime.now().strftime("%Y%m%d")
def start_str():
return (dt.datetime.now() - dt.timedelta(days=LOOKBACK)).strftime("%Y%m%d")
def read_parquet_safe(path):
try:
return pd.read_parquet(path)
except Exception: # noqa: BLE001
return None
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, dividend_type):
r = xd.get_market_data_ex([], [code], period="1d", start_time=start_str(),
end_time=today_str(), dividend_type=dividend_type)
return r.get(code) if r else None
def to_pdf(df):
return 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, # 手→股
})
def merge_write(code, dividend_type, kind):
"""返回 (added_rows, has_nan)。added=合并后比旧文件多的行数。"""
sym = code.split(".")[0]
prefix = prefix_of(sym)
df = fetch(code, dividend_type)
if df is None or not len(df):
return (0, False)
pdf = to_pdf(df)
has_nan = bool(pdf[["open", "high", "low", "close"]].isnull().any().any())
pdf["_y"] = pdf["date"].dt.year
added = 0
for yr, sub in pdf.groupby("_y"):
sub = sub.drop(columns=["_y"]).sort_values("date")
path = os.path.join(DATA, kind, str(int(yr)), f"{prefix}{sym}_daily.parquet")
old = read_parquet_safe(path) if os.path.exists(path) else None
old_n = len(old) if old is not None else 0
if old is not None and old_n:
merged = (pd.concat([old, sub])
.drop_duplicates("date", keep="last")
.sort_values("date"))
else:
merged = sub
added += max(len(merged) - old_n, 0)
atomic_write(path, merged)
return (added, has_nan)
def main():
end, start = today_str(), start_str()
log(f"start LOOKBACK={LOOKBACK} window={start}~{end}")
u = xd.get_stock_list_in_sector("沪深A股") or []
if not u:
log("FATAL: empty universeminiQMT 未连?)")
os._exit(2)
log(f"universe={len(u)}")
# 1. 批量下载增量到本地缓存(paced200/批,sleep 1s
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)
log("download phase done")
# 2. 逐只读 + 合并写
ok = fail = warn = 0
bars_added = 0
latest_dates = []
for i, code in enumerate(u):
try:
a1, w1 = merge_write(code, "none", "raw")
a2, w2 = merge_write(code, "front", "qfq")
bars_added += a1 + a2
if a1 + a2 > 0:
ok += 1
# 抽样记最新日期(前几只)
if len(latest_dates) < 5:
sym = code.split(".")[0]
prefix = prefix_of(sym)
yrd = dt.datetime.now().year
f = os.path.join(DATA, "raw", str(yrd), f"{prefix}{sym}_daily.parquet")
if os.path.exists(f):
d = read_parquet_safe(f)
if d is not None and len(d):
latest_dates.append((code, str(d["date"].max().date())))
if w1 or w2:
warn += 1
except Exception as e: # noqa: BLE001
fail += 1
if fail <= 5:
log(f"write err {code}: {e}")
if (i + 1) % 1000 == 0:
log(f"proc {i+1}/{len(u)} ok={ok} fail={fail} warn={warn} added={bars_added}")
log(f"DONE ok={ok} fail={fail} warn={warn} bars_added={bars_added}")
for c, d in latest_dates:
log(f"sample latest_date {c} -> {d}")
sys.stdout.flush()
os._exit(0 if fail == 0 else 1)
if __name__ == "__main__":
main()