b966fce1e9
- 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 引用
170 lines
5.7 KiB
Python
170 lines
5.7 KiB
Python
#!/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,volume(volume×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 universe(miniQMT 未连?)")
|
||
os._exit(2)
|
||
log(f"universe={len(u)}")
|
||
|
||
# 1. 批量下载增量到本地缓存(paced,200/批,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()
|