diff --git a/docs/data-platform/01-requirements.md b/docs/archive/data/01-requirements.md similarity index 100% rename from docs/data-platform/01-requirements.md rename to docs/archive/data/01-requirements.md diff --git a/docs/data-platform/02-p2-requirements.md b/docs/archive/data/02-p2-requirements.md similarity index 100% rename from docs/data-platform/02-p2-requirements.md rename to docs/archive/data/02-p2-requirements.md diff --git a/docs/data-platform/03-p3-requirements.md b/docs/archive/data/03-p3-requirements.md similarity index 100% rename from docs/data-platform/03-p3-requirements.md rename to docs/archive/data/03-p3-requirements.md diff --git a/docs/data/15min-data-design.md b/docs/archive/data/15min-data-design.md similarity index 100% rename from docs/data/15min-data-design.md rename to docs/archive/data/15min-data-design.md diff --git a/docs/superpowers/reports/2026-07-05-phase1-data-layer-completion.md b/docs/archive/data/2026-07-05-phase1-data-layer-completion.md similarity index 100% rename from docs/superpowers/reports/2026-07-05-phase1-data-layer-completion.md rename to docs/archive/data/2026-07-05-phase1-data-layer-completion.md diff --git a/docs/superpowers/plans/2026-07-05-plan1-data-layer.md b/docs/archive/data/2026-07-05-plan1-data-layer.md similarity index 100% rename from docs/superpowers/plans/2026-07-05-plan1-data-layer.md rename to docs/archive/data/2026-07-05-plan1-data-layer.md diff --git a/docs/superpowers/plans/2026-07-21-data-fusion-p0.md b/docs/archive/data/2026-07-21-data-fusion-p0.md similarity index 100% rename from docs/superpowers/plans/2026-07-21-data-fusion-p0.md rename to docs/archive/data/2026-07-21-data-fusion-p0.md diff --git a/docs/superpowers/plans/2026-07-22-data-arch-migration.md b/docs/archive/data/2026-07-22-data-arch-migration.md similarity index 100% rename from docs/superpowers/plans/2026-07-22-data-arch-migration.md rename to docs/archive/data/2026-07-22-data-arch-migration.md diff --git a/docs/superpowers/plans/2026-07-23-akshare-low-freq-schtask.md b/docs/archive/data/2026-07-23-akshare-low-freq-schtask.md similarity index 100% rename from docs/superpowers/plans/2026-07-23-akshare-low-freq-schtask.md rename to docs/archive/data/2026-07-23-akshare-low-freq-schtask.md diff --git a/docs/superpowers/plans/2026-07-23-csi1000-constituent-backfill.md b/docs/archive/data/2026-07-23-csi1000-constituent-backfill.md similarity index 100% rename from docs/superpowers/plans/2026-07-23-csi1000-constituent-backfill.md rename to docs/archive/data/2026-07-23-csi1000-constituent-backfill.md diff --git a/docs/superpowers/plans/2026-07-23-local-unified-provider.md b/docs/archive/data/2026-07-23-local-unified-provider.md similarity index 100% rename from docs/superpowers/plans/2026-07-23-local-unified-provider.md rename to docs/archive/data/2026-07-23-local-unified-provider.md diff --git a/docs/data-platform/daily-update-design.md b/docs/archive/data/daily-update-design.md similarity index 100% rename from docs/data-platform/daily-update-design.md rename to docs/archive/data/daily-update-design.md diff --git a/docs/portfolio_local_data_gaps.md b/docs/archive/data/portfolio_local_data_gaps.md similarity index 100% rename from docs/portfolio_local_data_gaps.md rename to docs/archive/data/portfolio_local_data_gaps.md diff --git a/docs/portfolio_local_unified_provider.md b/docs/archive/data/portfolio_local_unified_provider.md similarity index 100% rename from docs/portfolio_local_unified_provider.md rename to docs/archive/data/portfolio_local_unified_provider.md diff --git a/docs/static_data_cache_plan.md b/docs/archive/data/static_data_cache_plan.md similarity index 100% rename from docs/static_data_cache_plan.md rename to docs/archive/data/static_data_cache_plan.md diff --git a/docs/static_data_gaps_design.md b/docs/archive/data/static_data_gaps_design.md similarity index 100% rename from docs/static_data_gaps_design.md rename to docs/archive/data/static_data_gaps_design.md diff --git a/docs/data-platform/summary-p1.md b/docs/archive/data/summary-p1.md similarity index 100% rename from docs/data-platform/summary-p1.md rename to docs/archive/data/summary-p1.md diff --git a/docs/data-platform/README.md b/docs/data-platform/README.md new file mode 100644 index 0000000..a80e760 --- /dev/null +++ b/docs/data-platform/README.md @@ -0,0 +1,135 @@ +# 数据层总览(Data Layer) + +> A 股量化平台 **方案 A 数据层**单一权威记录。采集源 → 权威存储 → Provider → 策略全链路闭环。 +> 维护:数据 session。策略层接口对接见本文 §6;策略逻辑本身由策略 session 负责。 +> 深读设计依据:[`docs/superpowers/specs/2026-07-21-data-source-fusion-design.md`](../superpowers/specs/2026-07-21-data-source-fusion-design.md)(§14 权威层定稿)。 +> 历史中间设计/plan 已归档至 [`docs/archive/data/`](../archive/data/)。 + +--- + +## 1. 架构总览 + +``` +采集源(VPS定时任务) 权威存储层(VPS本地) 使用层(Provider) 策略层 +───────────────── ────────────────── ────────────────── ───────── +baostock ─┐ get_closes_panel 选股/轮动 +xtdata ─┼─► staging ─► 验证 ─► 合并 ─► dbbardata ─►┐ (BulletTrade +akshare ─┤ (validator) (merge) ├─ LocalUnifiedProvider 三策略) +sina/东财/腾讯─┘ constituent_unified─►┤ filters + valuation_baostock ─►┼─► get_fundamentals_df + 三表 parquet ─►────┘ get_limit_status_batch + bs_adjust_factor ... +``` + +**核心原则**:Provider 只读 VPS 本地数据,零 online;下载经 staging 隔离 → 验证 → 合并主库,绝不直接写主库。 + +--- + +## 2. 数据布局(VPS 本地,`C:\sanguo_vnpy_v2\data\`) + +| 存储 | 形态 | 内容 | 关键说明 | +|------|------|------|----------| +| `dbbardata`(quant_trading.db) | SQLite 表 | **唯一行情表**:日线(raw) + 15min | 含退市股 + ETF + 北交所920;schema `(id,symbol,exchange,datetime,interval,volume,turnover,open_interest,open/high/low/close_price)`;UNIQUE `(symbol,exchange,interval,datetime)`;WAL 模式 | +| `bs_adjust_factor` | SQLite 表 | 前复权因子 | `get_closes_panel(fq='qfq')` asof merge | +| `constituent_unified` | SQLite 表 | **20 指数成份股** | 9 宽基 + 000985 中证全指(~全市场) + 000928~000937 中证800十行业;`was_removed` 治幸存者偏差 | +| `valuation_baostock` | parquet/年 | pe / pb / isST | 2003-2026,按年 | +| 三表(akshare) | parquet | balance(221列)/income(170列)/cashflow(316列) | 财报季更新 | +| A股日线/分钟 | parquet | 兜底 | 回测 DB 为主,parquet 兜底 | + +**⚠️ 6 位码同名碰撞**:`000852/000905/000016/000985/000928~000937` 在 SZSE 是股票、在中证是指数点位。dbbardata 用 `exchange` 消歧:**`SSE` = 中证指数点位(约定),`SZSE` = 个股**。指数点位由 `sina_index_eod.py` 拉取入 `exchange=SSE`。 + +--- + +## 3. 采集源职责 + +| 源 | 职责 | 限制 | +|----|------|------| +| **baostock** | 个股日线(含退市)+ 15min + pe/pb | 单进程单登录,**不并发**(多连接→封 IP 6-24h);日 ≤ 48000 次 | +| **xtdata(miniQMT)** | ETF + 基金 + 北交所 920xxx | 沪深京全覆盖;volume 单位需 ×100 对齐 | +| **akshare** | 三表 + events + top10 股东 | 东财瞬时限流(断路器 + 夜间重试兜底) | +| **sina→东财→腾讯** | 14 指数点位 | 级联,腾讯兜底 5 个 sina 停更 | + +--- + +## 4. 增量管线(定时任务 schtask,VPS) + +| schtask | 时间 | 职责 | LOOKBACK | +|---------|------|------|----------| +| `sanguo-bs-eod` | 18:05 | baostock 个股日线 + 15min + pe/pb | 7 天 | +| `sanguo-idx-eod` | 18:30 | 14 指数点位(sina 级联) | — | +| `sanguo-ak-eod` | 19:00 | akshare 日线静态 | — | +| `sanguo-ak-events` | 19:30 | 事件 | — | +| `sanguo-xt-eod` | 21:00+ ⚠️ | ETF/基金/北交所920 | 30 天 | +| `sanguo-ak-stock` | 周六 | 个股全量(top10 等) | — | +| `sanguo-ak-quarter` | 财报季(APR,MAY,SEP,NOV) | 三表 | — | +| `sanguo-index` | 月度 16 号 19:50 | 成份股月度增量 | — | + +⚠️ `xt-eod` 原 18:40 与 `bs-eod` 18:05 写锁重叠(WAL 单写)→ 建议 21:00+ 错峰。 +`sanguo-index` STEP0:`parse_csindex_announce` 回溯公告治偏差(000852/932000)+ `--indices` 刷 11 新指数;STEP1-3 migrate→merge。 +`bs-eod` 已修 CLOSE_WAIT 卡死(per-stock commit + `_with_timeout` + 周期 relogin)。 + +--- + +## 5. 工作流与铁律 + +**下载链路(用户铁律)**: +``` +下载 → staging 隔离 → validator 验证(成功率95%,扣北交所) → 合并主库 → 推 NAS +``` +**绝不直接写主库**(下载质量不可控,曾出现覆盖截断整年)。 + +**约束**: +- baostock 单进程单登录不并发;日 ≤ 48000 次 +- 直连不走代理:`unset http_proxy https_proxy all_proxy` +- provider 读 VPS 本地,不调 online +- 数据层瑕疵报数据 session 根治,不在 provider 适配兜底 +- commit ≠ 部署 VPS:改完 `scp` 到 VPS + `findstr` 验证 +- VPS Windows:`python -X utf8`、反斜杠路径、GBK 控制台用 ASCII 脚本 +- Mac Mini 长任务前 `caffeinate -i -s` 防睡眠 + +--- + +## 6. Provider API(`LocalUnifiedProvider`) + +读本地零 online,治偏差(成份股并集 + 前视偏差修复)。14 个公开方法: + +| 方法 | 用途 | 备注 | +|------|------|------| +| `get_price(security, start, end, frequency, fields, count, fq)` | 单只行情 | 聚宽兼容 | +| **`get_closes_panel(symbols, start, end, interval='d', fq='raw')`** | 批量收盘价宽表 | `fq='qfq'` 批量前复权;UNION ALL 替 OR 链(340×);chunk=400;5128 只 33s | +| `get_index_stocks(index, date)` | 指数成份股 | 读 constituent_unified | +| `get_constituent(...)` | 成份股详情 | 治偏差 | +| **`get_fundamentals_df(stocks, date, fields=None)`** | 基本面 | `fields=` 短路(只算所需源表)+ ThreadPool;300 只 68s→7.4s(9.3×) | +| `get_value_metrics(security, date)` | 价值指标(单只) | — | +| **`get_value_metrics_batch(stocks, date)`** | 价值指标批量 | ThreadPool 包装 | +| `get_trade_days(start, end)` | 交易日历 | — | +| `get_security_info(security)` | 证券信息(单只) | — | +| **`get_security_info_batch(stocks)`** | 证券信息批量 | 2 SQL 替 N×2(ST/次新过滤提速) | +| **`get_limit_status_batch(codes, date)`** | 涨跌停/停牌批量 | 精确算 `high_limit=round(prev_close×(1+幅度),2)`;板块感知(主板10/创业科创20/北交30/ST5,历史ST from valuation_baostock.isST);停牌=volume==0 | +| `get_current_tick(security)` | tick | ⚠️ 无 last_price/high_limit 字段,filter 已改用 `get_limit_status_batch` | +| `get_split_dividend(security)` | 拆分分红/复权因子 | — | +| `get_all_securities(...)` | 全证券列表 | — | + +**四轮批量接口交付**(commit `f416a17`/`d2cd8fa`/`1cc9126`):行情 `get_closes_panel` / 财务 `get_fundamentals_df fields=` / filters+价值 `get_security_info_batch`+`get_value_metrics_batch` / 涨跌停停牌 `get_limit_status_batch`。 + +filters(`sanguo_portfolio/filters.py`):`filter_paused/limitup/limitdown` 已接入 `get_limit_status_batch`(向后兼容:无参=保留全部)。 + +--- + +## 7. 已知缺口与定论 + +| 项 | 状态 | 定论 | +|----|------|------| +| 行业治偏差(成份股历史调整) | 部分覆盖 | content HTML 解析已做(was_removed 8→99,主要覆盖 2009 + 零星 2016-2022);2010-2025 定期 csindex 无存档,**免费源穷尽**,用户接受残留偏差(不上 wind/choice) | +| 北交所 920xxx | ✅ 补全 | xtdata 独占(baostock/akshare 不覆盖),39 只,volume×100 对齐 | +| 实盘标的范围 | ✅ 定论 | 只做主板 + 创业板(`filter_kcbj_stock` 排除科创北交,用户未开户 50 万门槛);数据层仍全量补成份股治偏差,两层解耦 | +| 行业指数点位 | ✅ 修复 | 14 指数入 `exchange=SSE`,策略 03 不再永远判熊 | +| akshare 三表 | ✅ 鲁棒 | 原子写 + `--repair` + `is_parquet_healthy` | + +--- + +## 8. 待办(Phase 2,低优先) + +- 归档 15min 一次性灌库链(`backfill_15min_baostock` 等被 `tests/data/test_backfill_15min_hardening.py` import,需同步处理测试) +- 归档旧回填 import 链 + Mac `.sh` 链(`build_daily_from_xtdata`/`import_vnpy_*`/`raw_redownload` 等被 wrapper 引用,需 VPS `schtasks /query` 确认非活跃后归档) +- 行业治偏差 2010-2025(若有 wind/choice 授权再补) diff --git a/scripts/data_platform/_archive/README.md b/scripts/data_platform/_archive/README.md new file mode 100644 index 0000000..27a4ae1 --- /dev/null +++ b/scripts/data_platform/_archive/README.md @@ -0,0 +1,23 @@ +# 数据层归档脚本(_archive/) + +本目录存放**已完成使命的验证/诊断/一次性脚本**,不再参与日常增量管线。 +保留于 git 历史便于回溯;如需重跑,移动回 `scripts/data_platform/` 顶层即可。 + +## legacy/ — 探针 / 诊断 / 旧降级 / 一次性验证(2026-07-29 归档) + +| 脚本 | 类型 | 说明 | +|------|------|------| +| `probe_*.py`(12) | 探针 | 数据源/库/接口一次性 smoke 验证(akshare 状态/成份股/dbbardata 唯一性/退市/ETF/基本面/涨跌停/unified schema 等) | +| `dbbardata_probe.py` | 探针 | dbbardata 表结构与行数抽查 | +| `run_with_diag.py` / `diag_daily_update.ps1` | 诊断 | 带诊断输出的运行包装 | +| `test_mootdx_depth.py` / `test_baostock_daily_constituent_sample.py` | 一次性验证 | mootdx 深度 / baostock 日线成份股采样(非 tests/ 正式套件) | +| `resume_5yr_watcher.py` | 一次性 | 5 年全市场下载断点续传 watcher(已完成) | +| `fallback.py` / `realtime.py` | 旧降级 | 旧多源降级管理器(日线 akshare→腾讯 / 实时 新浪→东财→腾讯),方案 A 后由 bs_eod/xt_eod 接管 | + +归档前已验证:**零 import、无活跃 wrapper 引用**。 + +## backfill_15m/ — 15min 一次性灌库链(Phase 2 待归档) + +⚠️ 未归档。`backfill_15min_baostock` 被 `tests/data/test_backfill_15min_hardening.py` 正式 import, +`refresh_15min_daily` / `download_minute` / `download_15m_xtdata` / `raw_redownload` / `audit_data_layout` +存在交叉引用或 ops wrapper 依赖,需 Phase 2 评估后统一处理。 diff --git a/scripts/data_platform/dbbardata_probe.py b/scripts/data_platform/_archive/legacy/dbbardata_probe.py similarity index 100% rename from scripts/data_platform/dbbardata_probe.py rename to scripts/data_platform/_archive/legacy/dbbardata_probe.py diff --git a/scripts/data_platform/diag_daily_update.ps1 b/scripts/data_platform/_archive/legacy/diag_daily_update.ps1 similarity index 100% rename from scripts/data_platform/diag_daily_update.ps1 rename to scripts/data_platform/_archive/legacy/diag_daily_update.ps1 diff --git a/scripts/data_platform/fallback.py b/scripts/data_platform/_archive/legacy/fallback.py similarity index 100% rename from scripts/data_platform/fallback.py rename to scripts/data_platform/_archive/legacy/fallback.py diff --git a/scripts/data_platform/probe_akshare_status.py b/scripts/data_platform/_archive/legacy/probe_akshare_status.py similarity index 100% rename from scripts/data_platform/probe_akshare_status.py rename to scripts/data_platform/_archive/legacy/probe_akshare_status.py diff --git a/scripts/data_platform/probe_constituent.py b/scripts/data_platform/_archive/legacy/probe_constituent.py similarity index 100% rename from scripts/data_platform/probe_constituent.py rename to scripts/data_platform/_archive/legacy/probe_constituent.py diff --git a/scripts/data_platform/probe_dbbardata_unique.py b/scripts/data_platform/_archive/legacy/probe_dbbardata_unique.py similarity index 100% rename from scripts/data_platform/probe_dbbardata_unique.py rename to scripts/data_platform/_archive/legacy/probe_dbbardata_unique.py diff --git a/scripts/data_platform/probe_delisted_baostock.py b/scripts/data_platform/_archive/legacy/probe_delisted_baostock.py similarity index 100% rename from scripts/data_platform/probe_delisted_baostock.py rename to scripts/data_platform/_archive/legacy/probe_delisted_baostock.py diff --git a/scripts/data_platform/probe_delisted_ext.py b/scripts/data_platform/_archive/legacy/probe_delisted_ext.py similarity index 100% rename from scripts/data_platform/probe_delisted_ext.py rename to scripts/data_platform/_archive/legacy/probe_delisted_ext.py diff --git a/scripts/data_platform/probe_early_delisted_kline.py b/scripts/data_platform/_archive/legacy/probe_early_delisted_kline.py similarity index 100% rename from scripts/data_platform/probe_early_delisted_kline.py rename to scripts/data_platform/_archive/legacy/probe_early_delisted_kline.py diff --git a/scripts/data_platform/probe_etf.py b/scripts/data_platform/_archive/legacy/probe_etf.py similarity index 100% rename from scripts/data_platform/probe_etf.py rename to scripts/data_platform/_archive/legacy/probe_etf.py diff --git a/scripts/data_platform/probe_etf_v2.py b/scripts/data_platform/_archive/legacy/probe_etf_v2.py similarity index 100% rename from scripts/data_platform/probe_etf_v2.py rename to scripts/data_platform/_archive/legacy/probe_etf_v2.py diff --git a/scripts/data_platform/_archive/legacy/probe_fundamentals_panel.py b/scripts/data_platform/_archive/legacy/probe_fundamentals_panel.py new file mode 100644 index 0000000..420f6a5 --- /dev/null +++ b/scripts/data_platform/_archive/legacy/probe_fundamentals_panel.py @@ -0,0 +1,62 @@ +"""Smoke + timing: get_fundamentals_df fields= short-circuit + ThreadPool. + +ASCII-only (VPS GBK console safe). Run on VPS: + C:\\Python310\\python.exe -X utf8 probe_fundamentals_panel.py +""" +import os +import sys +import time + +os.environ.pop("http_proxy", None) +os.environ.pop("https_proxy", None) +os.environ.pop("all_proxy", None) + +sys.path.insert(0, r"C:\sanguo_vnpy_v2") + +from sanguo_portfolio.providers.local_unified_provider import LocalUnifiedProvider + +DB = r"C:\sanguo_vnpy_v2\data\quant_trading.db" +DATA = r"C:\sanguo_vnpy_v2\data" + +p = LocalUnifiedProvider({"db_path": DB, "data_dir": DATA}) + +# candidate pool: pull a few hundred codes from constituent_unified (000985 = full mkt) +try: + codes = p.get_index_stocks("000985", "2024-06-03") +except Exception as exc: + print("get_index_stocks failed:", exc) + codes = [] +codes = codes[:300] if codes else [] +jq = [c if "." in c else c + ".XSHE" for c in codes] +print("pool size:", len(jq)) +if not jq: + sys.exit(0) + +date = "2024-06-03" + +# warm caches once (first hit pays file open) to measure steady-ish state? No - +# measure COLD first-rebalance (the real pain): fields=None full read. +t0 = time.time() +df_none = p.get_fundamentals_df(jq, date=date) +t_none = time.time() - t0 + +# fresh provider to drop per-instance caches, measure fields= short-circuit cold +p2 = LocalUnifiedProvider({"db_path": DB, "data_dir": DATA}) +t0 = time.time() +df_fld = p2.get_fundamentals_df(jq, date=date, fields=["market_cap", "eps"]) +t_fld = time.time() - t0 + +print("fields=None : %6.2fs rows=%d cols=%d" % (t_none, len(df_none), len(df_none.columns))) +print("fields=[mkt,ep]: %6.2fs rows=%d cols=%d" % (t_fld, len(df_fld), len(df_fld.columns))) +if t_fld > 0: + print("speedup : %.1fx" % (t_none / t_fld)) + +# correctness: market_cap + eps match between the two +import pandas as pd +common = [c for c in df_fld.columns if c in df_none.columns] +for col in ("market_cap", "eps"): + a = df_none[col].reindex(df_fld.index) + b = df_fld[col] + mask = a.notna() & b.notna() + diff = (a[mask] - b[mask]).abs().max() if mask.any() else 0.0 + print("match %-12s: max_diff=%.6g" % (col, diff)) diff --git a/scripts/data_platform/_archive/legacy/probe_limit_status.py b/scripts/data_platform/_archive/legacy/probe_limit_status.py new file mode 100644 index 0000000..b4ef699 --- /dev/null +++ b/scripts/data_platform/_archive/legacy/probe_limit_status.py @@ -0,0 +1,54 @@ +"""Smoke: get_limit_status_batch on real dbbardata (window query + detection). + +ASCII-only (VPS GBK console safe). Run on VPS: + C:\\Python310\\python.exe -X utf8 probe_limit_status.py +""" +import collections +import os +import sys + +for k in ("http_proxy", "https_proxy", "all_proxy"): + os.environ.pop(k, None) + +sys.path.insert(0, r"C:\sanguo_vnpy_v2") +from sanguo_portfolio.providers.local_unified_provider import LocalUnifiedProvider + +DB = r"C:\sanguo_vnpy_v2\data\quant_trading.db" +DATA = r"C:\sanguo_vnpy_v2\data" +p = LocalUnifiedProvider({"db_path": DB, "data_dir": DATA}) + +DATE = "2024-06-03" +try: + codes = p.get_index_stocks("000985", DATE) +except Exception as exc: + print("get_index_stocks failed:", exc) + codes = [] +jq = [c if "." in c else c + ".XSHE" for c in codes[:800]] +print("pool:", len(jq), "date:", DATE) + +out = p.get_limit_status_batch(jq, DATE) +cnt = collections.Counter() +examples = {"up": [], "down": [], "paused": []} +for k, v in (out or {}).items(): + if v is None: + cnt["none"] += 1 + continue + if v["is_limit_up"]: + cnt["up"] += 1 + if len(examples["up"]) < 5: + examples["up"].append(k) + elif v["is_limit_down"]: + cnt["down"] += 1 + if len(examples["down"]) < 5: + examples["down"].append(k) + if v["is_paused"]: + cnt["paused"] += 1 + if len(examples["paused"]) < 5: + examples["paused"].append(k) + if not (v["is_limit_up"] or v["is_limit_down"] or v["is_paused"]): + cnt["normal"] += 1 + +print("counts:", dict(cnt)) +print("limit_up examples:", examples["up"]) +print("limit_down examples:", examples["down"]) +print("paused examples:", examples["paused"]) diff --git a/scripts/data_platform/_archive/legacy/probe_unified_aw.py b/scripts/data_platform/_archive/legacy/probe_unified_aw.py new file mode 100644 index 0000000..12eae15 --- /dev/null +++ b/scripts/data_platform/_archive/legacy/probe_unified_aw.py @@ -0,0 +1,134 @@ +# -*- coding: utf-8 -*- +"""UnifiedProvider + all_weather 数据链路诊断探针(VPS 跑, 快速版)。 + +每步带时间戳 + flush, 超时也能看卡哪。慢步骤降样本。 +定位: equity 重复日 / B_mean=0 / 数据缺失。 +""" +import sys +import os +import sqlite3 +import time +from collections import Counter + +try: + sys.stdout.reconfigure(line_buffering=True) +except Exception: + pass + +t0 = time.time() + + +def step(name): + print(f"\n=== {name} [+{time.time()-t0:.1f}s]", flush=True) + + +def line(k, v): + print(f"[{k}] {v}", flush=True) + + +VPS_ROOT = r"C:\sanguo_vnpy_v2" +DB = os.path.join(VPS_ROOT, "data", "quant_trading.db") + +os.environ.setdefault("DEFAULT_DATA_PROVIDER", "jqdata") +from unittest.mock import MagicMock +if "jqdatasdk" not in sys.modules: + _m = MagicMock() + _m.utils.assert_auth = lambda f: f + sys.modules["jqdatasdk"] = _m + +sys.path.insert(0, VPS_ROOT) + +step("STEP0 环境") +line("python", sys.version.split()[0]) +line("PKG", os.path.isdir(os.path.join(VPS_ROOT, "sanguo_portfolio"))) +line("DB", os.path.exists(DB)) +conn = sqlite3.connect(DB) +tabs = [r[0] for r in conn.execute("SELECT name FROM sqlite_master WHERE type='table'")] +line("tables", tabs) + +step("STEP0.5 混合 datetime 检测(单只抽样, 不全表 COUNT)") +# 单只 600519 抽样看格式(走索引, 快) +sample = conn.execute( + "SELECT datetime FROM dbbardata WHERE symbol='600519' AND exchange='SSE' " + "AND interval='d' ORDER BY datetime DESC LIMIT 5" +).fetchall() +line("600519 最近5条 datetime", [r[0] for r in sample]) +# DISTINCT 对比(单只, 索引内) +d_raw = conn.execute( + "SELECT COUNT(DISTINCT datetime) FROM dbbardata " + "WHERE symbol='600519' AND exchange='SSE' AND interval='d'" +).fetchone()[0] +d_sub = conn.execute( + "SELECT COUNT(DISTINCT substr(datetime,1,10)) FROM dbbardata " + "WHERE symbol='600519' AND exchange='SSE' AND interval='d'" +).fetchone()[0] +line("DISTINCT datetime(原始)", d_raw) +line("DISTINCT substr(datetime,1,10)(按日)", d_sub) +line("重复日数(原始-按日)", d_raw - d_sub) + +step("STEP1 get_trade_days 重复日期(equity_curve 重复 bug 根因)") +from sanguo_portfolio.providers import LocalUnifiedProvider +p = LocalUnifiedProvider({}) +days = p.get_trade_days(start_date="2024-01-02", end_date="2024-03-31") +strs = [str(d)[:10] for d in days] +line("trade_days total", len(days)) +line("unique dates", len(set(strs))) +dup = [d for d, c in Counter(strs).items() if c > 1] +line("DUP dates count", len(dup)) +line("DUP sample", dup[:5]) + +step("STEP2 成分股(constituent_unified 覆盖)") +for idx in ["000300", "399101", "399001", "000852"]: + try: + s = p.get_index_stocks(idx) + line(f"index_stocks {idx}", len(s)) + except Exception as e: + line(f"index_stocks {idx} ERR", repr(e)) + +step("STEP3 fundamentals 600519(单股, 关键字段)") +fdf = p.get_fundamentals_df(["600519.XSHG"], date="2024-03-29") +cols_chk = [ + "code", "market_cap", "circulating_market_cap", "pe_ratio", "pb_ratio", + "ps_ratio", "pcf_ratio", "eps", "roe", "roa", "gross_profit_margin", + "net_profit_margin", "inc_revenue_year_on_year", "roic", +] +for c in cols_chk: + if c in fdf.columns: + line(f" {c}", fdf[c].iloc[0]) + else: + line(f" {c}", "MISSING_COL") + +step("STEP4 _trend_mean 小样本复算(hs300 前40, B_mean=0 根因)") +import numpy as np +hs300 = p.get_index_stocks("000300") +line("hs300 size", len(hs300)) +# 只取前 40 只做 fundamentals(提速), top20 by circ_mktcap +sample40 = hs300[:40] +fdf2 = p.get_fundamentals_df(sample40, date="2024-03-29") +line("fdf2 shape", fdf2.shape) +if "circulating_market_cap" in fdf2.columns: + line("circ_mktcap nonNaN", int(fdf2["circulating_market_cap"].notna().sum())) +fdf2s = fdf2.sort_values("circulating_market_cap", ascending=False, na_position="last") +blst = list(fdf2s.index)[:20] +line("blst(20)", blst) +df = p.get_price(blst, end_date="2024-03-29", frequency="1d", fields=["close"], count=10, panel=False) +line("trend get_price isNone", df is None) +if df is not None: + line("trend get_price shape", df.shape) + line("trend cols", list(df.columns)) + line("time dtype", df["time"].dtype if "time" in df.columns else "NO_TIME") + print(df.head(3).to_string(), flush=True) + try: + pivot = df.pivot(index="time", columns="code", values="close") + line("pivot shape", pivot.shape) + if len(pivot) >= 2: + change = (pivot.iloc[-1] / pivot.iloc[0] - 1) * 100 + arr = np.nan_to_num(change.to_numpy()) + line("B_mean manual", float(np.mean(arr))) + line("change nonZero count", int((arr != 0).sum())) + else: + line("pivot rows<2", len(pivot)) + except Exception as e: + line("pivot ERR", repr(e)) + +step("DONE") diff --git a/scripts/data_platform/probe_unified_schema.py b/scripts/data_platform/_archive/legacy/probe_unified_schema.py similarity index 100% rename from scripts/data_platform/probe_unified_schema.py rename to scripts/data_platform/_archive/legacy/probe_unified_schema.py diff --git a/scripts/data_platform/realtime.py b/scripts/data_platform/_archive/legacy/realtime.py similarity index 100% rename from scripts/data_platform/realtime.py rename to scripts/data_platform/_archive/legacy/realtime.py diff --git a/scripts/data_platform/resume_5yr_watcher.py b/scripts/data_platform/_archive/legacy/resume_5yr_watcher.py similarity index 100% rename from scripts/data_platform/resume_5yr_watcher.py rename to scripts/data_platform/_archive/legacy/resume_5yr_watcher.py diff --git a/scripts/data_platform/_archive/legacy/run_with_diag.py b/scripts/data_platform/_archive/legacy/run_with_diag.py new file mode 100644 index 0000000..98064eb --- /dev/null +++ b/scripts/data_platform/_archive/legacy/run_with_diag.py @@ -0,0 +1,28 @@ +# -*- coding: utf-8 -*- +"""monkey-patch engine 关键方法加诊断, 跑回测看 ETF cancel 根因(不改源码)。""" +import sys, os, logging +os.environ.setdefault("DEFAULT_DATA_PROVIDER", "jqdata") +from unittest.mock import MagicMock +m = MagicMock(); m.utils.assert_auth = lambda f: f +sys.modules.setdefault("jqdatasdk", m) +logging.basicConfig(level=logging.WARNING, format="%(message)s") + +from bullet_trade.core import engine as eng + +_orig_calc = eng.BacktestEngine._calculate_order_amount +def calc(self, order, cp): + r = _orig_calc(self, order, cp) + print(f"[ENG_DIAG] {order.security} cp={cp} amount={r} tgt_val={getattr(order,'_target_value',None)} is_tgt={getattr(order,'_is_target_value',None)} order_amt={getattr(order,'amount',None)}", flush=True) + return r +eng.BacktestEngine._calculate_order_amount = calc + +_orig_bp = eng.BacktestEngine._resolve_base_exec_price +def bp(self, security, current_dt, fq_mode): + r = _orig_bp(self, security, current_dt, fq_mode) + print(f"[ENG_DIAG_BP] {security} dt={current_dt} fq={fq_mode} -> {r}", flush=True) + return r +eng.BacktestEngine._resolve_base_exec_price = bp + +sys.argv = ['runner', '--provider', 'unified', '--max-pool', '30', '--start', '2024-01-02', '--end', '2024-01-31', '--cash', '1000000'] +from sanguo_portfolio.runner_backtest import main +main() diff --git a/scripts/data_platform/test_baostock_daily_constituent_sample.py b/scripts/data_platform/_archive/legacy/test_baostock_daily_constituent_sample.py similarity index 100% rename from scripts/data_platform/test_baostock_daily_constituent_sample.py rename to scripts/data_platform/_archive/legacy/test_baostock_daily_constituent_sample.py diff --git a/scripts/data_platform/test_mootdx_depth.py b/scripts/data_platform/_archive/legacy/test_mootdx_depth.py similarity index 100% rename from scripts/data_platform/test_mootdx_depth.py rename to scripts/data_platform/_archive/legacy/test_mootdx_depth.py