#!/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 = list(set( (xd.get_stock_list_in_sector("沪深A股") or []) + (xd.get_stock_list_in_sector("沪深ETF") or []) + (xd.get_stock_list_in_sector("沪深基金") 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()