diff --git a/scripts/data_platform/daily_update_xtdata.py b/scripts/data_platform/daily_update_xtdata.py new file mode 100644 index 0000000..ac1f274 --- /dev/null +++ b/scripts/data_platform/daily_update_xtdata.py @@ -0,0 +1,169 @@ +#!/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() diff --git a/scripts/data_platform/import_vnpy_daily_fast.py b/scripts/data_platform/import_vnpy_daily_fast.py index f97caab..2a5bd16 100644 --- a/scripts/data_platform/import_vnpy_daily_fast.py +++ b/scripts/data_platform/import_vnpy_daily_fast.py @@ -12,8 +12,8 @@ import sys import time from pathlib import Path -DB_PATH = '/tmp/quant_trading_import.db' -DAILY_DIR = '/Volumes/stock/A股数据/日线数据/daily/' +DB_PATH = os.environ.get('VNPY_DB_PATH', '/tmp/quant_trading_import.db') +DAILY_DIR = os.environ.get('DAILY_DIR', '/Volumes/stock/A股数据/日线数据/daily/') def parse_filename(filename): @@ -41,7 +41,9 @@ def import_year(conn, year): if code is None: continue try: - df = pd.read_parquet(f, columns=['date', 'open', 'high', 'low', 'close', 'volume', 'amount']) + df = pd.read_parquet(f, columns=['date', 'open', 'high', 'low', 'close', 'volume']) + if 'amount' not in df.columns: + df['amount'] = 0.0 if df.empty: continue df['symbol'] = code