#!/usr/bin/env python3 # -*- coding: utf-8 -*- """sina_index_eod.py — sanguo-idx-eod (18:30 schtask): 中证指数日线点位 -> dbbardata('d', SSE)。 背景: dbbardata 6 位码同时是中证指数码与 SZSE 真实股票码 (000928~000937/000852/000905/ 000016/000985), 历史只灌了 SZSE 股票行, 中证指数点位从未入库。000300 无碰撞, 已有 SSE 指数行 (close 几千)。本脚本把 14 个中证指数点位灌入 exchange='SSE', 与 SZSE 股票行 (exchange='SZSE') 隔离 — provider 的 .XSHG->SSE 映射本就如此。 数据源 (实证 2026-07-28, 级联取第一个"最近 30 天内有数据"的源): 主 sina: ak.stock_zh_index_daily(symbol="sh000928") cols=[date(str), open, high, low, close, volume] 9/14 良好; 000929/000930/000936/000937/000985 sina 停在 2016-06-13 (源残缺)。 兜底1 东财: ak.index_zh_a_hist(symbol, period="daily") — VPS IP 持续 RemoteDisconnected。 兜底2 腾讯: ak.stock_zh_index_daily_tx(symbol="sh000929") — schema 不同! cols=[date(date 对象), open, close, high, low, amount] (注意: 无 volume, amount 是成交额, 列顺序 close 在 high 前)。对全 14 都有数据, 是当前兜底主力。 字段映射: volume 取 sina/东财 的 volume, 腾讯 volume=0; turnover=腾讯 amount, 其他源 0。 策略只用 close, volume/turnover 不影响; OHLC 三源都有。 源选择 (surgical): 不把腾讯改 primary 是为了保留 9 个 sina 良好 code 的 volume 数据 (腾讯无 volume 列, 切 primary 会把已入库的 volume 清零, 是回归)。腾讯仅在 sina stale 时兜底。 幂等: dbbardata UNIQUE(symbol,exchange,datetime,interval) -> INSERT OR REPLACE 自动幂等。 只写 exchange='SSE' 行, 绝不动 exchange='SZSE' 股票行。 硬约束 (用户铁律): 顶部清 http_proxy/https_proxy/all_proxy (直连不走代理); 单线程, 每次 akshare 调用 sleep 1.2s; 禁止并发; 腾讯 tqdm 进度条不阻塞 (函数会跑完)。 用法: python sina_index_eod.py [--lookback-days N] [--dry-run] # N 默认 30 (env LOOKBACK_DAYS) python sina_index_eod.py --lookback-days 3650 # 首次回灌 """ import argparse import os import sys import time # 硬约束: 直连不走代理 (akshare sina/eastmoney 都是国内源, 走代理反而挂) for _k in ("http_proxy", "https_proxy", "all_proxy", "HTTP_PROXY", "HTTPS_PROXY", "ALL_PROXY"): os.environ.pop(_k, None) import sqlite3 import pandas as pd # 同目录 import (与 xt_eod/bs_eod 一致, schtask 工作目录 = scripts/data_platform) sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from dbbardata_utils import normalize_daily_dt # noqa: E402 DB = r"C:\sanguo_vnpy_v2\data\quant_trading.db" T0 = time.time() # 15 个中证指数 (与 SZSE 股票码碰撞, 仅灌 SSE 点位行) # 000938 = 中证 1000 等权 (策略需求: RPS/动量基准对比, 加 2026-07-28) CODES = [ "000928", "000929", "000930", "000931", "000932", "000933", "000934", "000935", "000936", "000937", "000938", "000852", "000905", "000016", "000985", ] SINA_SLEEP = 1.2 # 单线程限速: 每次 akshare 调用后 sleep (避封 IP) FALLBACK_RECENT_DAYS = 30 # sina 最后日期 < today-30 才触发东财兜底 def log(m): print(f"[IDX-EOD {time.time()-T0:.0f}s] {m}", flush=True) # 统一 schema: 三源都对齐到这 7 列 (turnover 缺则 0; volume 腾讯缺则 0) _UNIFIED_COLS = ["date", "open", "high", "low", "close", "volume", "turnover"] def _normalize_df(df, col_map, default_turnover=0.0, default_volume=0.0): """重命名 + 补默认列 -> _UNIFIED_COLS。缺失列用 default 兜底。""" if df is None or len(df) == 0: return pd.DataFrame() df2 = df.rename(columns=col_map) if "turnover" not in df2.columns: df2["turnover"] = default_turnover if "volume" not in df2.columns: df2["volume"] = default_volume for col in ("date", "open", "high", "low", "close"): if col not in df2.columns: return pd.DataFrame() # 必需列缺失 # date 归一为 'YYYY-MM-DD' 字符串 (腾讯 date 是 date 对象, str() 即得) df2["date"] = df2["date"].astype(str).str[:10] return df2[_UNIFIED_COLS] def fetch_sina(code): """sina 主源 -> _UNIFIED_COLS。返空 df 表示 sina 无数据。""" import akshare as ak df = ak.stock_zh_index_daily(symbol=f"sh{code}") time.sleep(SINA_SLEEP) return _normalize_df( df, col_map={"date": "date", "open": "open", "high": "high", "low": "low", "close": "close", "volume": "volume"}, default_turnover=0.0, ) def fetch_eastmoney(code): """东财兜底 -> _UNIFIED_COLS。VPS IP 持续 RemoteDisconnected 时返空。""" import akshare as ak try: df = ak.index_zh_a_hist(symbol=code, period="daily") except (ConnectionError, OSError, Exception) as e: log(f" {code} eastmoney err: {type(e).__name__}: {str(e)[:80]}") return pd.DataFrame() time.sleep(SINA_SLEEP) return _normalize_df( df, col_map={"日期": "date", "开盘": "open", "最高": "high", "最低": "low", "收盘": "close", "成交量": "volume", "成交额": "turnover"}, default_turnover=0.0, ) def fetch_tencent(code): """腾讯兜底 -> _UNIFIED_COLS。 实证 schema: cols=[date(date 对象), open, close, high, low, amount] — 注意: - date 是 datetime.date 对象, str() 后取 [:10] 得 'YYYY-MM-DD' - 无 volume 列, 用 amount (成交额) 当 turnover; volume=0 - 列顺序 close 在 high 之前 (col_map 重命名兼容) """ import akshare as ak try: df = ak.stock_zh_index_daily_tx(symbol=f"sh{code}") except (ConnectionError, OSError, Exception) as e: log(f" {code} tencent err: {type(e).__name__}: {str(e)[:80]}") return pd.DataFrame() time.sleep(SINA_SLEEP) return _normalize_df( df, col_map={"date": "date", "open": "open", "close": "close", "high": "high", "low": "low", "amount": "turnover"}, default_volume=0.0, ) def fetch_with_fallback(code, today): """级联 sina → 东财 → 腾讯, 取第一个"最后日期 >= today-30"的源。 Returns: (df, source, warn_msg) - df: _UNIFIED_COLS DataFrame (可能为空) - source: 'sina'|'eastmoney'|'tencent'|'none' - warn_msg: 兜底/降级原因 (sina 直接命中时为 None) 全部源都 stale 时, 用最新的 stale df (有总比没有好, log WARN)。 """ candidates = ( ("sina", fetch_sina), ("eastmoney", fetch_eastmoney), ("tencent", fetch_tencent), ) stale_best = None # (df, source, last_dt, reason) warn_parts = [] for name, fn in candidates: try: df = fn(code) except Exception as e: warn_parts.append(f"{name} 异常 {type(e).__name__}") continue if df is None or len(df) == 0: warn_parts.append(f"{name} 空") continue try: last_dt = pd.Timestamp(df["date"].iloc[-1]).normalize() except Exception: last_dt = pd.Timestamp.min if (today - last_dt).days <= FALLBACK_RECENT_DAYS: # 命中: 新鲜数据 warn = None if not warn_parts else ("上游失效: " + "; ".join(warn_parts) + f" → {name} 命中") return df, name, warn # stale: 保留作最后兜底 if stale_best is None or last_dt > stale_best[2]: stale_best = (df, name, last_dt, f"{name} 停 {last_dt.date()}") warn_parts.append(f"{name} 停 {last_dt.date()}") if stale_best is not None: df, src, _, reason = stale_best return df, src, "全源 stale (>30d); " + "; ".join(warn_parts) return pd.DataFrame(), "none", "全源失败: " + "; ".join(warn_parts) def to_dbbardata_rows(code, df): """_UNIFIED_COLS DataFrame -> dbbardata 行元组列表 (symbol,SSE,d,...)。""" if df is None or len(df) == 0: return [] out = [] for row in df.itertuples(index=False): dt_str = normalize_daily_dt(str(row.date)) if not dt_str or dt_str == "None": continue try: out.append(( code, "SSE", dt_str, "d", float(row.volume), float(row.turnover), 0.0, float(row.open), float(row.high), float(row.low), float(row.close), )) except (ValueError, TypeError) as e: log(f" {code} skip row {dt_str}: {e}") return out def main(): ap = argparse.ArgumentParser() ap.add_argument("--lookback-days", type=int, default=int(os.environ.get("LOOKBACK_DAYS", "30"))) ap.add_argument("--dry-run", action="store_true") ap.add_argument("--codes", type=str, default="", help="逗号分隔覆盖默认 14 个 code (调试用)") args = ap.parse_args() codes = (args.codes.split(",") if args.codes else CODES) codes = [c.strip() for c in codes if c.strip()] today = pd.Timestamp.now().normalize() start_ts = today - pd.Timedelta(days=args.lookback_days) start_str = start_ts.strftime("%Y-%m-%d") log(f"start={start_str} lookback={args.lookback_days}d codes={len(codes)}" f"{' [DRY-RUN]' if args.dry_run else ''}") conn = sqlite3.connect(DB, timeout=60) conn.execute("PRAGMA busy_timeout = 60000") conn.execute("PRAGMA journal_mode = WAL") ok_codes = 0 total_rows = 0 src_counts = {"sina": 0, "eastmoney": 0, "tencent": 0, "none": 0} code_stats = [] # (code, source, rows, min_close, max_close, max_dt, warn?) conn.execute("BEGIN") try: for code in codes: df, source, warn_msg = fetch_with_fallback(code, today) src_counts[source] = src_counts.get(source, 0) + 1 # 过滤 lookback 范围 if not df.empty: df = df.copy() df["date"] = pd.to_datetime(df["date"], format="mixed").dt.strftime("%Y-%m-%d") df = df[df["date"] >= start_str] rows = to_dbbardata_rows(code, df) if not rows: log(f" {code} {source} 无符合 lookback 行" + (f" ({warn_msg})" if warn_msg else "")) code_stats.append((code, source, 0, None, None, None, warn_msg or "lookback 范围内无行")) continue 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 (?,?,?,?,?,?,?,?,?,?,?)", rows, ) ok_codes += 1 total_rows += len(rows) closes = [r[10] for r in rows] max_dt = max(r[2] for r in rows) min_c = round(min(closes), 2) max_c = round(max(closes), 2) code_stats.append((code, source, len(rows), min_c, max_c, max_dt, warn_msg)) log(f" {code} {source} rows={len(rows)} close=[{min_c},{max_c}] last={max_dt}" + (f" WARN: {warn_msg}" if warn_msg else "")) conn.execute("COMMIT") except Exception as e: conn.execute("ROLLBACK") log(f"FATAL rollback: {type(e).__name__}: {e}") conn.close() return 1 conn.close() src_summary = ", ".join(f"{k}={v}" for k, v in src_counts.items() if v) log(f"DONE codes_ok={ok_codes}/{len(codes)} rows={total_rows} sources[{src_summary}]" f"{' [DRY-RUN]' if args.dry_run else ''}") return 0 if __name__ == "__main__": raise SystemExit(main())