#!/usr/bin/env python3 # -*- coding: utf-8 -*- """xt_eod.py — sanguo-xt-eod (方案A schtask 18:40): xtata ETF/基金/北交所 EOD 增量 -> dbbardata('d')。 baostock 只取 type=1 股票, 不覆盖 ETF/基金/北交所个股 -> xtata 独占 (spec §14)。 北交所个股 baostock 不覆盖, 此处补; 沪深个股仍由 bs_eod 灌, 避免 dbbardata 两源冲突。 - universe = 沪深ETF ∪ 沪深基金 ∪ 北交所920xxx (中证2000 成份股, baostock 不覆盖, xtata 独占) - download_history_data2 批量 paced -> 本地缓存 - get_market_data_ex raw(dividend_type=none) -> dbbardata('d') INSERT OR REPLACE - volume 手->股 (×100, 与 daily_update_xtdata 同口径) - 无限流, 单进程 download 不并发 用法: python xt_eod.py [--limit N] [--dry-run] # 日常增量 (LOOKBACK=30 天) python xt_eod.py --full-bj # 北交所 backfill (start=20240101) """ import argparse import datetime as dt import sqlite3 import time try: from xtquant import xtdata as xd except ImportError: xd = None # mac 单测 exc_of 时 xd=None, VPS 跑 main() 会 return 2 import pandas as pd from dbbardata_utils import normalize_daily_dt DB = r"C:\sanguo_vnpy_v2\data\quant_trading.db" LOOKBACK = int(__import__("os").environ.get("LOOKBACK_DAYS", "30")) T0 = time.time() def log(m): print(f"[XT-EOD {time.time()-T0:.0f}s] {m}", flush=True) def prefix_of(sym): return "sh" if sym[:2] in ("51", "56", "58", "50") else ("sh" if sym[:2] == "60" else "sz") def exc_of(sym): # 920 是 3 位前缀 (北交所), 必须在 2 位 SSE/SZSE 判断前优先, 否则 sym[:2]='92' 落 SZSE if sym[:3] == "920": return "BJSE" return "SSE" if sym[:2] in ("51", "56", "58", "50", "60", "68") else "SZSE" def main(): ap = argparse.ArgumentParser() ap.add_argument("--limit", type=int, default=0) ap.add_argument("--dry-run", action="store_true") ap.add_argument("--full-bj", action="store_true", help="北交所 920xxx backfill: start=20240101 (默认与 ETF 同 LOOKBACK)") args = ap.parse_args() if xd is None: log("FATAL xtquant 未装(VPS-only)") return 2 end = dt.datetime.now().strftime("%Y%m%d") etf_start = (dt.datetime.now() - dt.timedelta(days=LOOKBACK)).strftime("%Y%m%d") bj_start = "20240101" if args.full_bj else etf_start log(f"start window ETF/基金={etf_start} 北交所={bj_start}~{end} (full_bj={args.full_bj})") # 沪深 ETF/基金 (xtata 独占, baostock 不覆盖) etf_codes = list(set( (xd.get_stock_list_in_sector("沪深ETF") or []) + (xd.get_stock_list_in_sector("沪深基金") or []) )) # 北交所 920xxx (中证2000 成份股, baostock 不覆盖, xtata 独占) bj_codes = [] try: _c = sqlite3.connect(DB, timeout=30) bj_raw = [r[0] for r in _c.execute( "SELECT DISTINCT code FROM constituent_unified " "WHERE index_code='932000' AND code LIKE '920%'" )] _c.close() bj_codes = [f"{c}.BJ" for c in bj_raw] except Exception as e: log(f"WARN constituent_unified 920 read err: {e}") log(f"universe ETF/基金={len(etf_codes)} 北交所={len(bj_codes)}") u = etf_codes + bj_codes if not u: log("FATAL empty universe (miniQMT 未连?)") return 2 if args.limit: u = u[:args.limit] def _start_of(code): return bj_start if code.split(".")[0].startswith("920") else etf_start # download paced: 按 start 分组避免 download_history_data2 单 start 限制 BATCH = 200 for st in ({etf_start, bj_start}): sub = [c for c in u if _start_of(c) == st] if not sub: continue for i in range(0, len(sub), BATCH): try: xd.download_history_data2(sub[i:i+BATCH], "1d", st, end, lambda d, p: None) except Exception as e: log(f"dl @{st} @{i} err: {e}") time.sleep(1.0) log("download done") conn = sqlite3.connect(DB, timeout=60) conn.execute("PRAGMA busy_timeout = 60000") conn.execute("PRAGMA journal_mode = WAL") ok = fail = empty = rows = 0 conn.execute("BEGIN") try: for i, code in enumerate(u): sym = code.split(".")[0] try: r = xd.get_market_data_ex([], [code], period="1d", start_time=_start_of(code), end_time=end, dividend_type="none") df = r.get(code) if r else None if df is None or not len(df): empty += 1 continue db = pd.DataFrame({ "symbol": sym, "exchange": exc_of(sym), # datetime 归一纯日期 (dbbardata 双行根治方案A) "datetime": [normalize_daily_dt( f"{str(idx)[:4]}-{str(idx)[4:6]}-{str(idx)[6:8]}") for idx in df.index], "interval": "d", "volume": (df["volume"].astype(float).values * 100), "turnover": df["amount"].astype(float).values, "open_interest": 0.0, "open_price": df["open"].astype(float).values, "high_price": df["high"].astype(float).values, "low_price": df["low"].astype(float).values, "close_price": df["close"].astype(float).values, }) if not args.dry_run: conn.executemany( "INSERT OR REPLACE INTO dbbardata " "(symbol,exchange,datetime,interval,volume,turnover,open_interest," "open_price,high_price,low_price,close_price) VALUES (?,?,?,?,?,?,?,?,?,?,?)", db.itertuples(index=False, name=None)) rows += len(db) ok += 1 except Exception as e: fail += 1 if fail <= 5: log(f"{code} err: {e}") if (i+1) % 200 == 0: log(f"进度 {i+1}/{len(u)} ok={ok} empty={empty} fail={fail} rows={rows}") conn.execute("COMMIT") except Exception as e: conn.execute("ROLLBACK") log(f"FATAL rollback: {e}") conn.close() return 1 conn.close() log(f"DONE ok={ok} empty={empty} fail={fail} rows={rows}" f"{' [DRY-RUN]' if args.dry_run else ''}") return 0 if __name__ == "__main__": raise SystemExit(main())