From 4966d30a115673b4826eb761ed846d074b245eae Mon Sep 17 00:00:00 2001 From: claude_dev Date: Wed, 2 Sep 2026 10:43:11 +0800 Subject: [PATCH] =?UTF-8?q?feat(provider):=20A=E6=A1=A3=E6=88=AA=E9=9D=A2?= =?UTF-8?q?=E6=95=B0=E6=8D=AE=E4=BD=BF=E7=94=A8=E5=B1=82=E5=87=BA=E5=8F=A3?= =?UTF-8?q?=E2=80=94=E2=80=94=E6=B3=9B=E5=9E=8B=E5=BA=95=E5=BA=A7+?= =?UTF-8?q?=E6=B6=A8=E5=81=9C=E6=B1=A0=E9=97=A8=E9=9D=A2=E5=8F=8C=E5=B1=82?= =?UTF-8?q?=E6=8E=A5=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 09-02任务①落地。§15采集的14类panel数据此前零读取出口,策略/因子session 想用只能手搓parquet路径;经业界对照(dsa get_limit_up_pool/tushare limit_list_d U/D/Z参数化/聚宽per-domain/OpenBB TET)定稿「泛型底座+域 语义门面」: ①get_event_panel(event_type,date/start+end,trading_days_only=True): 保真底座,白名单14类统一通道(=采集注册表口径,hot_rank未挂载不入列), 返回akshare原样中文列+trade_date(文件名合成);缺文件=合法缺失→空df; trading_days_only默认用get_trade_days滤非交易日——治快照族「周中法定 假日文件=上一交易日态复制品」的区间双计数(采集层只按周一~五落盘, 周六根本无文件,风险面=周中假日);=False按周一~五枚举(=采集口径) ②get_limit_pool(kind=zt|zbgc|dtgc,...):涨停池门面(温度计/炸板率/行业 集中度消费契约),kind三合一←tushare limit_list_d;英文标准列 code/consecutive_boards/seal_amount/break_count/industry…(键名对齐dsa); 映射表基于VPS真实parquet实测(dtgc实测含动态市盈率/封单资金/板上成交额/ 连续跌停/开板次数,与文档口径有差);akshare改中文列名时本层吸收漂移, 有行缺源列→DataSchemaError fail-fast;真空日(dtgc 0跌停空文件无列) →仍返标准列空表schema稳定 ③实现=TET Fetcher(fetchers/event_panel.py双Fetcher同文件,price.py先例), FETCHERS注册表+event_panel/limit_pool两键(未来MCP出口零成本);base.py 新增_DateRangeParams跨字段校验(date/start+end互斥二选一);LocalUnified Provider(回测)与SanguoMiniQmtProvider(实盘,委托self._unified)双侧同款 ——方法面钉死测试强制无缺口,双侧同schema支持副本对照;miniQMT/QMT/ PTrade均无此类接口(业界惯例=外部源补),实盘同读本地文件 测试:21新测试(白名单fail-fast/日期语义/保真读/假日滤/门面schema/ 列漂移fail-fast/真空日空表/注册表)+全量793绿;活文档§16+§15指引行 [vps] --- .../2026-07-21-data-source-fusion-design.md | 36 ++- .../providers/fetchers/__init__.py | 9 + sanguo_portfolio/providers/fetchers/base.py | 71 ++++- .../providers/fetchers/event_panel.py | 204 +++++++++++++ .../providers/local_unified_provider.py | 48 ++++ .../providers/sanguo_fundamentals.py | 35 +++ tests/portfolio/test_event_panel_fetchers.py | 268 ++++++++++++++++++ 7 files changed, 669 insertions(+), 2 deletions(-) create mode 100644 sanguo_portfolio/providers/fetchers/event_panel.py create mode 100644 tests/portfolio/test_event_panel_fetchers.py diff --git a/docs/superpowers/specs/2026-07-21-data-source-fusion-design.md b/docs/superpowers/specs/2026-07-21-data-source-fusion-design.md index ba0476a..47d0bf1 100644 --- a/docs/superpowers/specs/2026-07-21-data-source-fusion-design.md +++ b/docs/superpowers/specs/2026-07-21-data-source-fusion-design.md @@ -3,7 +3,7 @@ > 日期:2026-07-21 创建 | 基于 brainstorming + 4 源全能力调查(akshare / baostock / miniQMT / csindex) > **状态:活文档(唯一数据设计文档)** —— 2026-09-02 起用户钦定:本文件是数据架构唯一权威设计, > 每次数据架构变更必须同批更新到最新状态;演进采用「日期节追加,新节覆盖旧节相关决策」模式。 -> 最后更新:2026-09-02(A 档截面数据域全量上生产,§15) +> 最后更新:2026-09-02(A 档截面数据域全量上生产,§15;使用层 provider 接口,§16) --- @@ -296,6 +296,8 @@ VPS config 的 daily_dir/raw_dir/qfq_dir/minute_15_dir 改 `C:\sanguo_vnpy_v2\da ## 15. 增补 2026-09-02:A 档截面数据域全量上生产(本节为最新,覆盖前文相关决策) +> 本节数据的**使用层读取接口见 §16**(get_event_panel / get_limit_pool,同日落地)。 + > 背景:调研 [daily_stock_analysis](https://github.com/ZhuLinsen/daily_stock_analysis)(Gitea 镜像 `sanguo/daily_stock_analysis`) > 的数据采集面,双机实验(Mac + VPS 生产 IP 49.232.102.198)后落地。四 commit > db789d6→020d7cd→00cfd45→46a0569,CI#303 NAS 验证绿 + vps-deploy#304 绿,tag=46a0569。 @@ -352,6 +354,38 @@ baostock 栈=bs_* 系列,xtdata 栈=daily_update_xtdata,新浪指数=sina_index_ =免费当日热点行业分布信号,涨停行业集中度原料)/换手率/市值等 - 温度计三支柱=涨跌停(hot_rank 待批)+雪球;情绪周期/打板/行业轮动因子原料自 2026-08 起逐日积累 +## 16. 增补 2026-09-02:使用层 provider 接口(A 档数据的读取出口) + +> §15 采集落盘的数据在本节获得 provider 出口。设计经业界对照定稿(2026-09-02 用户拍板): +> dsa `get_limit_up_pool` / tushare `limit_list_d`(U/D/Z 参数化)/ 聚宽 per-domain 方法 / OpenBB TET +> ——共识三条:**接口轴=数据域**、**列名在数据层边界标准化**(无一家透传原始中文列)、 +> **泛型传输+语义入口分层**(tushare 底层也是 `pro.query(api_name, params)` 泛型通道)。 +> miniQMT/QMT/PTrade 均无涨停池/龙虎榜类接口(业界惯例=外部源补),实盘 provider 同样读本地文件。 + +### 16.1 两层出口(「泛型底座 + 域语义门面」) + +| 接口 | 定位 | 返回列 | +|---|---|---| +| `get_event_panel(event_type, date / start+end, trading_days_only=True)` | **保真底座**:白名单 14 类统一通道(§15.2 全部+gdhs),研究/因子探索用 | akshare 原样中文列 + `trade_date`(YYYY-MM-DD,文件名合成);**别把底座列名当生产契约** | +| `get_limit_pool(kind="zt"\|"zbgc"\|"dtgc", date / start+end, ...)` | **涨停池门面**:kind 三合一(← tushare limit_list_d),唯一有消费契约的域先建(温度计/炸板率/行业集中度) | 英文标准列:`code/consecutive_boards/seal_amount/break_count/first_seal_time/industry/limit_stat`…(键名对齐 dsa) | + +- 实现=TET Fetcher(`fetchers/event_panel.py`,FETCHERS 注册表新增 `event_panel`/`limit_pool`), + LocalUnifiedProvider(回测)与 SanguoMiniQmtProvider(实盘,委托 `self._unified`)双侧同款 + ——方法面钉死测试(issue #35)强制无缺口;双侧同 schema 支持副本对照验证 +- **列名漂移吸收**:akshare 改中文列名时,门面映射表(`_LIMIT_POOL_COLUMNS`,基于 VPS 真实 + parquet 实测 2026-09-02,dtgc 实测含 动态市盈率/封单资金/板上成交额/连续跌停/开板次数)吸收, + 下游英文键不炸;有行但缺源列→DataSchemaError fail-fast +- **trading_days_only=True(默认)**:用 `get_trade_days` 滤非交易日——治快照族「周中法定假日文件 + =上一交易日态复制品」的区间双计数;`=False` 按周一~五枚举(=采集侧口径) +- 缺文件(洞/未来日/从未采集)=合法缺失→空 DataFrame;真空日(dtgc 0 跌停)→门面仍返标准列空表 +- hot_rank 未挂载不在白名单(拍板后与采集注册表、PANEL_TYPES 三处同步加) +- 新数据类型接入=底座白名单加一行(零新方法);某域有消费契约后再升格门面(YAGNI,不加投机门面) + +### 16.2 MCP 出口预留 + +FETCHERS 注册表新增两键即自动可被未来 provider MCP server 暴露 +(`Fetcher.fetch(ctx=provider, **kwargs)` 直接调用,fetchers 包设计初衷)。 + ## 参考(调查来源) - xtdata 官方:https://dict.thinktrader.net/nativeApi/xtdata.html diff --git a/sanguo_portfolio/providers/fetchers/__init__.py b/sanguo_portfolio/providers/fetchers/__init__.py index 6fdfbbd..64896d8 100644 --- a/sanguo_portfolio/providers/fetchers/__init__.py +++ b/sanguo_portfolio/providers/fetchers/__init__.py @@ -10,12 +10,15 @@ Phase 1(本包)只加新接口不动老接口——老接口照常可用;Phase 2 from .base import ( ConstituentQueryParams, DataSchemaError, + EventPanelQueryParams, FundamentalsQueryParams, + LimitPoolQueryParams, PanelQueryParams, PriceQueryParams, validate_df_schema, ) from .constituent import ConstituentFetcher +from .event_panel import EventPanelFetcher, LimitPoolFetcher from .fundamentals import FundamentalsFetcher from .price import PanelFetcher, PriceFetcher @@ -25,6 +28,8 @@ FETCHERS = { "closes_panel": PanelFetcher, "constituent": ConstituentFetcher, "fundamentals": FundamentalsFetcher, + "event_panel": EventPanelFetcher, + "limit_pool": LimitPoolFetcher, } __all__ = [ @@ -32,11 +37,15 @@ __all__ = [ "PanelFetcher", "ConstituentFetcher", "FundamentalsFetcher", + "EventPanelFetcher", + "LimitPoolFetcher", "FETCHERS", "PriceQueryParams", "PanelQueryParams", "ConstituentQueryParams", "FundamentalsQueryParams", + "EventPanelQueryParams", + "LimitPoolQueryParams", "DataSchemaError", "validate_df_schema", ] diff --git a/sanguo_portfolio/providers/fetchers/base.py b/sanguo_portfolio/providers/fetchers/base.py index fa4fb29..ba2e4da 100644 --- a/sanguo_portfolio/providers/fetchers/base.py +++ b/sanguo_portfolio/providers/fetchers/base.py @@ -24,7 +24,7 @@ from __future__ import annotations from datetime import datetime from typing import Any, Optional, Sequence, Union -from pydantic import BaseModel, ConfigDict, field_validator +from pydantic import BaseModel, ConfigDict, field_validator, model_validator # fq 合法值(与老接口 get_price/get_closes_panel 取值一致) FQ_VALUES = ("raw", "qfq", "pre", "前复权") @@ -210,6 +210,75 @@ class FundamentalsQueryParams(_QueryBase): return fs +# 事件/快照 panel 白名单(2026-09-02 A 档使用层出口):与采集注册表 +# (akshare_static_download.py)及 static_vintage_check.PANEL_TYPES 口径一致; +# hot_rank 墙测量未拍板挂载,故不在列(拍板后三处同步加)。 +EVENT_PANEL_TYPES = frozenset({ + "dragon_tiger", "block_trade", "margin_sse", "restricted", + "zt_pool", "zt_pool_zbgc", "zt_pool_dtgc", + "fund_flow_industry", "fund_flow_concept", + "ths_industry", "ths_concept", + "xueqiu_hot", "sina_sector", "gdhs", +}) + + +class _DateRangeParams(_QueryBase): + """date / start+end 二选一语义(事件 panel 族共用)。 + + - 单日: ``date``;区间: ``start``+``end``(含端点,二缺一/互斥/倒序 → 报错) + """ + + date: Optional[str] = None + start: Optional[str] = None + end: Optional[str] = None + + @field_validator("date", "start", "end") + @classmethod + def _v_dates(cls, v): + return _norm_date(v) + + @model_validator(mode="after") + def _check_date_semantics(self): + if self.date and (self.start or self.end): + raise ValueError("date 与 start/end 互斥(单日用 date,区间用 start+end)") + if not self.date and not self.start: + raise ValueError("date 与 start 必须给其一") + if (self.start and not self.end) or (self.end and not self.start): + raise ValueError("区间查询必须同时给 start 和 end") + if self.start and self.end and self.start > self.end: + raise ValueError("start 不能晚于 end") + return self + + +class EventPanelQueryParams(_DateRangeParams): + """get_event_panel 入参(泛型保真通道)。""" + + event_type: str + trading_days_only: bool = True + + @field_validator("event_type") + @classmethod + def _v_event_type(cls, v): + if v not in EVENT_PANEL_TYPES: + raise ValueError( + f"event_type={v!r} 不在白名单(合法: {sorted(EVENT_PANEL_TYPES)})") + return v + + +class LimitPoolQueryParams(_DateRangeParams): + """get_limit_pool 入参(涨停池门面,kind 三合一 ← tushare limit_list_d U/D/Z)。""" + + kind: str = "zt" + trading_days_only: bool = True + + @field_validator("kind") + @classmethod + def _v_kind(cls, v): + if v not in ("zt", "zbgc", "dtgc"): + raise ValueError("kind 必须是 'zt'|'zbgc'|'dtgc'(涨停|炸板|跌停)") + return v + + def validate_df_schema( df: Any, *, diff --git a/sanguo_portfolio/providers/fetchers/event_panel.py b/sanguo_portfolio/providers/fetchers/event_panel.py new file mode 100644 index 0000000..8379866 --- /dev/null +++ b/sanguo_portfolio/providers/fetchers/event_panel.py @@ -0,0 +1,204 @@ +"""事件/快照 panel Fetcher: get_event_panel(保真底座) + get_limit_pool(涨停池门面)。 + +设计(2026-09-02,活文档 §16):「泛型底座 + 域语义门面」——OpenBB TET 架构下的 +两层出口,业界共识(dsa/tushare limit_list_d/聚宽)全走「接口轴=数据域 + 列名在 +数据层边界标准化」: + +- **EventPanelFetcher(保真底座)**: 任意白名单类型统一通道,返回 akshare 原样 + 中文列 + ``trade_date``(文件名合成)。定位=研究/因子探索通道(extract 原样 + 输出,不做 transform 标准化)——**别把底座列名当生产契约**,契约走门面。 +- **LimitPoolFetcher(涨停池门面)**: kind 三合一(zt/zbgc/dtgc,对齐 tushare + ``limit_list_d`` 的 limit_type=U/D/Z 参数化),中文列→英文标准键(键名对齐 + dsa ``get_limit_up_pool``)。akshare 改中文列名时,本层映射吸收漂移、 + 下游代码不炸;有行但缺 代码 → DataSchemaError fail-fast。 + +数据布局(采集侧 §15): ``{data_dir}/static/{type}/{YYYYMMDD}_{type}.parquet``; +per-date 语义=每工作日 1 文件(节假日空/快照族=上一交易日态复制品)。 + +关键正确性——``trading_days_only=True``(默认)用 ``ctx.get_trade_days`` 滤掉 +非交易日: 快照族节假日文件是上一交易日的复制品,区间聚合会双计数(§15 语义); +zt_pool 族节假日是真空文件,滤掉无损。raw 模式(=False)按周一~周五枚举,语义 +对齐采集侧 ``build_per_date_units``。 + +列名映射基于 VPS 真实 parquet 实测(2026-09-02): dtgc 实测含 动态市盈率/ +封单资金/板上成交额/连续跌停/开板次数——与 akshare 文档口径有差,以实测为准。 +""" +from __future__ import annotations + +from datetime import datetime, timedelta +from pathlib import Path +from typing import TYPE_CHECKING, List, Optional + +import pandas as pd + +from .base import ( + DataSchemaError, + EventPanelQueryParams, + LimitPoolQueryParams, + validate_df_schema, +) + +if TYPE_CHECKING: + from ..local_unified_provider import LocalUnifiedProvider + +# kind → 采集侧类型名(=data/static 子目录名) +_LIMIT_POOL_TYPES = { + "zt": "zt_pool", + "zbgc": "zt_pool_zbgc", + "dtgc": "zt_pool_dtgc", +} + +# 中文列 → 英文标准键(实测 2026-09-02;键序=输出列序;序号=源行号无信息量,弃) +_LIMIT_POOL_COLUMNS = { + "zt": { + "代码": "code", "名称": "name", "涨跌幅": "change_pct", "最新价": "price", + "成交额": "amount", "流通市值": "float_market_cap", "总市值": "total_market_cap", + "换手率": "turnover_rate", "封板资金": "seal_amount", + "首次封板时间": "first_seal_time", "最后封板时间": "last_seal_time", + "炸板次数": "break_count", "涨停统计": "limit_stat", + "连板数": "consecutive_boards", "所属行业": "industry", + }, + "zbgc": { + "代码": "code", "名称": "name", "涨跌幅": "change_pct", "最新价": "price", + "涨停价": "limit_price", "成交额": "amount", "流通市值": "float_market_cap", + "总市值": "total_market_cap", "换手率": "turnover_rate", "涨速": "surge_speed", + "首次封板时间": "first_seal_time", "炸板次数": "break_count", + "涨停统计": "limit_stat", "振幅": "amplitude", "所属行业": "industry", + }, + "dtgc": { + "代码": "code", "名称": "name", "涨跌幅": "change_pct", "最新价": "price", + "成交额": "amount", "流通市值": "float_market_cap", "总市值": "total_market_cap", + "动态市盈率": "pe_ttm", "换手率": "turnover_rate", "封单资金": "seal_amount", + "最后封板时间": "last_seal_time", "板上成交额": "board_amount", + "连续跌停": "consecutive_limit_downs", "开板次数": "open_count", + "所属行业": "industry", + }, +} + + +def _workdays(start: str, end: str) -> List[str]: + """周一~周五枚举(=False 模式,语义对齐采集侧 build_per_date_units)。""" + out: List[str] = [] + cur = datetime.strptime(start, "%Y-%m-%d") + last = datetime.strptime(end, "%Y-%m-%d") + while cur <= last: + if cur.weekday() < 5: + out.append(cur.strftime("%Y%m%d")) + cur += timedelta(days=1) + return out + + +def _resolve_query_dates( + query, ctx: "LocalUnifiedProvider" +) -> List[str]: + """date/start+end → 待读日期列表(YYYYMMDD)。""" + if query.date: + return [query.date.replace("-", "")] + if query.trading_days_only: + days = ctx.get_trade_days( + start_date=query.start, end_date=query.end) + return [d.strftime("%Y%m%d") for d in days] + return _workdays(query.start, query.end) + + +def _read_type_frames( + ctx: "LocalUnifiedProvider", event_type: str, query +) -> pd.DataFrame: + """读 ``static/{type}/{YYYYMMDD}_{type}.parquet`` 并 concat(唯一 IO)。 + + - 文件缺失(洞/未来日/从未采集)= 合法缺失,跳过(空 DataFrame 兜底) + - 真空文件(0 行)贡献零行,不占列 + - ``trade_date``(YYYY-MM-DD)由文件名合成追加在列尾 + """ + type_dir = Path(ctx.data_dir) / "static" / event_type + frames: List[pd.DataFrame] = [] + for d in _resolve_query_dates(query, ctx): + p = type_dir / f"{d}_{event_type}.parquet" + if not p.exists(): + continue + df = pd.read_parquet(p) + if not len(df): + continue + df = df.copy() + df["trade_date"] = f"{d[:4]}-{d[4:6]}-{d[6:]}" + frames.append(df) + if not frames: + return pd.DataFrame() + return pd.concat(frames, ignore_index=True) + + +class EventPanelFetcher: + """``get_event_panel``: 泛型保真通道(akshare 原样中文列 + trade_date)。""" + + @staticmethod + def transform_query(**kwargs) -> EventPanelQueryParams: + return EventPanelQueryParams(**kwargs) + + @staticmethod + def extract_data( + query: EventPanelQueryParams, ctx: "LocalUnifiedProvider" + ) -> pd.DataFrame: + return _read_type_frames(ctx, query.event_type, query) + + @staticmethod + def transform_data( + query: EventPanelQueryParams, + ctx: "LocalUnifiedProvider", + df: pd.DataFrame, + ) -> pd.DataFrame: + """底座=保真通道:零清洗零重命名(空 df=合法缺失,无 schema 钉死)。""" + return df + + @classmethod + def fetch(cls, ctx: "LocalUnifiedProvider", **kwargs) -> pd.DataFrame: + query = cls.transform_query(**kwargs) + df = cls.extract_data(query, ctx) + return cls.transform_data(query, ctx, df) + + +class LimitPoolFetcher: + """``get_limit_pool``: 涨停池门面(kind 三合一,英文标准列)。 + + transform_data=OpenBB「标准数据模型」时刻:中文→英文映射 + ``code`` + 列 fail-fast;真空日仍返标准列空表(schema 稳定,消费方免判空列)。 + """ + + @staticmethod + def transform_query(**kwargs) -> LimitPoolQueryParams: + return LimitPoolQueryParams(**kwargs) + + @staticmethod + def extract_data( + query: LimitPoolQueryParams, ctx: "LocalUnifiedProvider" + ) -> pd.DataFrame: + return _read_type_frames( + ctx, _LIMIT_POOL_TYPES[query.kind], query) + + @staticmethod + def transform_data( + query: LimitPoolQueryParams, + ctx: "LocalUnifiedProvider", + df: pd.DataFrame, + ) -> pd.DataFrame: + mapping = _LIMIT_POOL_COLUMNS[query.kind] + standard = list(mapping.values()) + ["trade_date"] + if not len(df): + return pd.DataFrame(columns=standard) + missing = [cn for cn in mapping if cn not in df.columns] + if missing: + raise DataSchemaError( + f"limit_pool[{query.kind}]: 源列缺失 {missing}" + f"(akshare 列名漂移?df.columns={list(df.columns)})") + out = df[list(mapping)].rename(columns=mapping).copy() + out["trade_date"] = df["trade_date"].values + validate_df_schema( + out, required=["code"], + context=f"get_limit_pool(kind={query.kind})", + ) + return out[standard] + + @classmethod + def fetch(cls, ctx: "LocalUnifiedProvider", **kwargs) -> pd.DataFrame: + query = cls.transform_query(**kwargs) + df = cls.extract_data(query, ctx) + return cls.transform_data(query, ctx, df) diff --git a/sanguo_portfolio/providers/local_unified_provider.py b/sanguo_portfolio/providers/local_unified_provider.py index eb0a9e2..2040f70 100644 --- a/sanguo_portfolio/providers/local_unified_provider.py +++ b/sanguo_portfolio/providers/local_unified_provider.py @@ -966,3 +966,51 @@ class LocalUnifiedProvider(DataProvider): # type: ignore[misc] """TET 版 get_fundamentals_df(多股财务/估值)。Phase 3 起与老接口同一实现(strict 契约)。""" from .fetchers.fundamentals import FundamentalsFetcher return FundamentalsFetcher.fetch(self, stocks=stocks, date=date, fields=fields) + + # ==================== 事件/快照 panel(A 档使用层出口,2026-09-02) ==================== + def get_event_panel( + self, + event_type: str, + date: Optional[Union[str, datetime]] = None, + start: Optional[Union[str, datetime]] = None, + end: Optional[Union[str, datetime]] = None, + trading_days_only: bool = True, + ) -> pd.DataFrame: + """泛型保真通道: 读 ``static/{event_type}/{YYYYMMDD}_{type}.parquet``。 + + 白名单 14 类(与采集注册表/缺日检查口径一致,详见活文档 §16): + 涨停池×3/资金流×2/同花顺×2/雪球/新浪行业/龙虎榜/大宗/两融/解禁/gdhs。 + + - 返回 akshare **原样中文列** + ``trade_date``(YYYY-MM-DD,文件名合成)—— + 保真探索通道,**别把列名当生产契约**(契约走 ``get_limit_pool`` 门面) + - ``trading_days_only=True``(默认)滤非交易日:快照族节假日文件=上一交易 + 日态复制品,区间聚合会双计数(§15 语义);``False``=周一~五枚举(采集口径) + - 文件缺失(洞/未来日)= 合法缺失 → 空 DataFrame + """ + from .fetchers.event_panel import EventPanelFetcher + return EventPanelFetcher.fetch( + self, event_type=event_type, date=date, start=start, end=end, + trading_days_only=trading_days_only) + + def get_limit_pool( + self, + kind: str = "zt", + date: Optional[Union[str, datetime]] = None, + start: Optional[Union[str, datetime]] = None, + end: Optional[Union[str, datetime]] = None, + trading_days_only: bool = True, + ) -> pd.DataFrame: + """涨停池门面(kind 三合一 ← tushare ``limit_list_d`` 的 U/D/Z 参数化)。 + + kind: ``"zt"``(涨停池) | ``"zbgc"``(炸板池) | ``"dtgc"``(跌停池)。 + + - 英文标准列(``code/consecutive_boards/seal_amount/break_count/ + industry``…,键名对齐 dsa);akshare 改中文列名时本层映射吸收漂移 + - 有行但缺源列(如 代码)→ DataSchemaError fail-fast + - 真空日(如 dtgc 0 跌停)→ 标准列空表(schema 稳定) + - 炸板率口径(消费契约): 炸板池/(涨停池+炸板池),按行数在策略层算 + """ + from .fetchers.event_panel import LimitPoolFetcher + return LimitPoolFetcher.fetch( + self, kind=kind, date=date, start=start, end=end, + trading_days_only=trading_days_only) diff --git a/sanguo_portfolio/providers/sanguo_fundamentals.py b/sanguo_portfolio/providers/sanguo_fundamentals.py index a0f20de..95ba518 100644 --- a/sanguo_portfolio/providers/sanguo_fundamentals.py +++ b/sanguo_portfolio/providers/sanguo_fundamentals.py @@ -415,6 +415,41 @@ class SanguoMiniQmtProvider(MiniQMTProvider): # type: ignore[misc] logger.warning("get_value_metrics_batch 委托本地库失败,返空: %s", exc) return {} + def get_event_panel( + self, + event_type: str, + date: Optional[Union[str, datetime]] = None, + start: Optional[Union[str, datetime]] = None, + end: Optional[Union[str, datetime]] = None, + trading_days_only: bool = True, + ) -> pd.DataFrame: + """事件/快照 panel 保真通道:委托本地 unified(与回测同源同口径)。 + + 白名单/契约见 ``LocalUnifiedProvider.get_event_panel``(活文档 §16)。 + miniQMT 无涨停池/龙虎榜类接口(业界共识=外部源补),实盘读昨晚 19:30 + ak-events 落盘的本地文件。 + """ + return self._unified.get_event_panel( + event_type, date=date, start=start, end=end, + trading_days_only=trading_days_only) + + def get_limit_pool( + self, + kind: str = "zt", + date: Optional[Union[str, datetime]] = None, + start: Optional[Union[str, datetime]] = None, + end: Optional[Union[str, datetime]] = None, + trading_days_only: bool = True, + ) -> pd.DataFrame: + """涨停池门面(kind 三合一):委托本地 unified(与回测同源同口径)。 + + 英文标准列契约见 ``LocalUnifiedProvider.get_limit_pool``;实盘/回测 + 双侧同 schema,支持副本对照验证。 + """ + return self._unified.get_limit_pool( + kind, date=date, start=start, end=end, + trading_days_only=trading_days_only) + def get_price_ex( self, security: Union[str, List[str]], diff --git a/tests/portfolio/test_event_panel_fetchers.py b/tests/portfolio/test_event_panel_fetchers.py new file mode 100644 index 0000000..a38c69d --- /dev/null +++ b/tests/portfolio/test_event_panel_fetchers.py @@ -0,0 +1,268 @@ +# -*- coding: utf-8 -*- +"""事件/快照 panel 数据 Fetcher 测试(2026-09-02 A 档使用层出口)。 + +两层接口(设计=泛型底座 + 涨停池语义门面,对齐 OpenBB TET/tushare limit_list_d): +1. ``EventPanelFetcher``(get_event_panel): 保真通道——akshare 原样中文列 + + ``trade_date``(文件名合成);新数据类型零成本接入;缺文件=合法空。 +2. ``LimitPoolFetcher``(get_limit_pool): 涨停池门面——kind 三合一 + (zt/zbgc/dtgc ← tushare limit_list_d 的 U/D/Z 参数化),英文标准列 + (键名对齐 dsa get_limit_up_pool),列缺失 fail-fast(吸收 akshare 列漂移)。 + +关键契约: +- ``trading_days_only=True``(默认)用 get_trade_days 滤非交易日——治快照族 + 「节假日文件=上一交易日态复制品」的区间双计数(§15 语义)。 +- 列名映射基于 VPS 真实 parquet 实测(2026-09-02),非 akshare 文档臆测: + dtgc 实测含 动态市盈率/封单资金/板上成交额/连续跌停/开板次数。 +- 真空日(如 dtgc 0 跌停)空文件无列:门面仍返标准列空表(schema 稳定)。 +""" +from __future__ import annotations + +import sqlite3 + +import pandas as pd +import pytest +from pydantic import ValidationError + +from sanguo_portfolio.providers.fetchers import ( + FETCHERS, + DataSchemaError, + EventPanelFetcher, + LimitPoolFetcher, +) +from sanguo_portfolio.providers.local_unified_provider import LocalUnifiedProvider + +# 真实列名(VPS 20260819_zt_pool_dtgc.parquet 实测) +ZT_COLS = { + "序号", "代码", "名称", "涨跌幅", "最新价", "成交额", "流通市值", "总市值", + "换手率", "封板资金", "首次封板时间", "最后封板时间", "炸板次数", + "涨停统计", "连板数", "所属行业", +} +ZBGC_COLS = { + "序号", "代码", "名称", "涨跌幅", "最新价", "涨停价", "成交额", "流通市值", + "总市值", "换手率", "涨速", "首次封板时间", "炸板次数", "涨停统计", + "振幅", "所属行业", +} +DTGC_COLS = { + "序号", "代码", "名称", "涨跌幅", "最新价", "成交额", "流通市值", "总市值", + "动态市盈率", "换手率", "封单资金", "最后封板时间", "板上成交额", + "连续跌停", "开板次数", "所属行业", +} + +# 门面标准键(每 kind 的完整 schema;trade_date 由读取层追加) +ZT_KEYS = { + "code", "name", "change_pct", "price", "amount", "float_market_cap", + "total_market_cap", "turnover_rate", "seal_amount", "first_seal_time", + "last_seal_time", "break_count", "limit_stat", "consecutive_boards", + "industry", "trade_date", +} +ZBGC_KEYS = ZT_KEYS - {"seal_amount", "last_seal_time", "consecutive_boards"} | { + "limit_price", "surge_speed", "amplitude"} +DTGC_KEYS = ZT_KEYS - {"first_seal_time", "break_count", + "consecutive_boards", "limit_stat"} | { + "pe_ttm", "board_amount", "consecutive_limit_downs", "open_count"} + + +def _mk_zt_row(code: str, boards: int, industry: str = "白酒") -> dict: + return { + "序号": 1, "代码": code, "名称": "样本股", "涨跌幅": 10.0, "最新价": 11.0, + "成交额": 1e8, "流通市值": 5e9, "总市值": 8e9, "换手率": 5.0, + "封板资金": 2e8, "首次封板时间": "092500", "最后封板时间": "150000", + "炸板次数": 0, "涨停统计": "2天/2板", "连板数": boards, "所属行业": industry, + } + + +@pytest.fixture +def provider(tmp_path): + """交易日历(2024-06-18/20,19=周三非交易日) + static 各型样本文件。 + + 2024-06-19=周中「法定假日」形态: 采集层按周一~五落盘会为它写 + 「上一交易日态复制品」(快照族语义)——用它钉死 trading_days_only 过滤。 + """ + db = tmp_path / "quant_trading.db" + c = sqlite3.connect(str(db)) + c.execute( + "CREATE TABLE dbbardata(symbol TEXT, exchange TEXT, datetime TEXT, " + "interval TEXT, volume REAL, turnover REAL, open_interest REAL, " + "open_price REAL, high_price REAL, low_price REAL, close_price REAL)" + ) + for d in ("2024-06-18", "2024-06-20"): + c.execute( + "INSERT INTO dbbardata VALUES(?,?,?,?,?,?,?,?,?,?,?)", + ("600519", "SSE", f"{d} 00:00:00", "d", 1000, 1e6, 0, + 1.0, 1.0, 1.0, 1.0)) + c.commit() + c.close() + + static = tmp_path / "static" + + def _write(t: str, day: str, rows: list) -> None: + d = static / t + d.mkdir(parents=True, exist_ok=True) + pd.DataFrame(rows, columns=list(rows[0]) if rows else None).to_parquet( + d / f"{day}_{t}.parquet", index=False) + + _write("zt_pool", "20240618", [_mk_zt_row("600519", 2), _mk_zt_row("000001", 3)]) + _write("zt_pool", "20240619", [_mk_zt_row("600519", 1)]) # 假日复制品 + _write("zt_pool_zbgc", "20240619", [{ + "序号": 1, "代码": "300999", "名称": "炸板股", "涨跌幅": 9.5, "最新价": 8.8, + "涨停价": 9.9, "成交额": 5e7, "流通市值": 2e9, "总市值": 3e9, "换手率": 12.0, + "涨速": 3.2, "首次封板时间": "093000", "炸板次数": 2, "涨停统计": "1天/0板", + "振幅": 8.0, "所属行业": "半导体"}]) + _write("zt_pool_dtgc", "20240618", []) # 真空日(0 跌停)=空文件无列 + _write("gdhs", "20240331", [ + {"代码": "600519", "股东户数": 80000, "区间增减": -0.05}, + {"代码": "000001", "股东户数": 500000, "区间增减": 0.02}]) + return LocalUnifiedProvider({"db_path": str(db), "data_dir": str(tmp_path)}) + + +# ======================== QueryParams strict(fail-fast) ======================== + +class TestEventPanelQueryParams: + + def test_unknown_type_rejected(self, provider): + with pytest.raises(ValidationError, match="no_such_type"): + provider.get_event_panel("no_such_type", date="2024-06-18") + + def test_hot_rank_not_mounted(self, provider): + """hot_rank 墙测量未拍板挂载——白名单拒绝(与 PANEL_TYPES 口径一致)。""" + with pytest.raises(ValidationError, match="hot_rank"): + provider.get_event_panel("hot_rank", date="2024-06-18") + + def test_date_and_start_exclusive(self, provider): + with pytest.raises(ValidationError, match="互斥"): + provider.get_event_panel( + "zt_pool", date="2024-06-18", start="2024-06-18", + end="2024-06-20") + + def test_neither_date_nor_start(self, provider): + with pytest.raises(ValidationError, match="其一"): + provider.get_event_panel("zt_pool") + + def test_start_requires_end(self, provider): + with pytest.raises(ValidationError, match="start 和 end"): + provider.get_event_panel("zt_pool", start="2024-06-18") + + def test_bad_date_format(self, provider): + with pytest.raises(ValidationError): + provider.get_event_panel("zt_pool", date="20240618") + + +class TestLimitPoolQueryParams: + + def test_unknown_kind_rejected(self, provider): + with pytest.raises(ValidationError, match="kind"): + provider.get_limit_pool(kind="ztbad", date="2024-06-18") + + def test_date_semantics_shared(self, provider): + with pytest.raises(ValidationError, match="其一"): + provider.get_limit_pool(kind="zt") + + +# ======================== 泛型底座(保真通道) ======================== + +class TestEventPanelFetch: + + def test_single_date_raw_columns_plus_trade_date(self, provider): + df = provider.get_event_panel("zt_pool", date="2024-06-18") + assert len(df) == 2 + assert set(df.columns) >= ZT_COLS + assert "trade_date" in df.columns + assert set(df["trade_date"]) == {"2024-06-18"} + + def test_range_trading_days_only_drops_holiday_file(self, provider): + """默认过滤: 周中假日的复制品文件必须被滤掉(防区间双计数)。""" + df = provider.get_event_panel( + "zt_pool", start="2024-06-18", end="2024-06-20") + # 18(2行);19=非交易日复制品被滤;20 无文件(合法洞) + assert len(df) == 2 + assert set(df["trade_date"]) == {"2024-06-18"} + + def test_range_raw_mode_includes_holiday(self, provider): + df = provider.get_event_panel( + "zt_pool", start="2024-06-18", end="2024-06-20", + trading_days_only=False) + assert len(df) == 3 + assert "2024-06-19" in set(df["trade_date"]) + + def test_missing_dir_returns_empty(self, provider): + """从未采集的类型(xueqiu_hot 目录不存在)=合法缺失,空 df 不报错。""" + df = provider.get_event_panel("xueqiu_hot", date="2024-06-18") + assert isinstance(df, pd.DataFrame) and df.empty + + def test_empty_file_day_contributes_nothing(self, provider): + df = provider.get_event_panel( + "zt_pool_dtgc", start="2024-06-18", end="2024-06-20") + assert df.empty # 18=真空文件,19/20 无文件 + + def test_gdhs_period_read(self, provider): + """gdhs per-period: unit={季度末}_gdhs,同一底座直接可读。""" + df = provider.get_event_panel("gdhs", date="2024-03-31") + assert len(df) == 2 + assert set(df["trade_date"]) == {"2024-03-31"} + + +# ======================== 涨停池门面(标准模型) ======================== + +class TestLimitPoolFetch: + + def test_zt_english_schema(self, provider): + df = provider.get_limit_pool(kind="zt", date="2024-06-18") + assert set(df.columns) == ZT_KEYS + row = df.iloc[0] + assert row["consecutive_boards"] == 2 + assert row["seal_amount"] == 2e8 + assert row["break_count"] == 0 + assert row["industry"] == "白酒" + assert row["limit_stat"] == "2天/2板" + + def test_zbgc_schema(self, provider): + df = provider.get_limit_pool(kind="zbgc", date="2024-06-19") + assert set(df.columns) == ZBGC_KEYS + assert df.iloc[0]["break_count"] == 2 + assert df.iloc[0]["surge_speed"] == 3.2 + + def test_dtgc_empty_day_keeps_standard_schema(self, provider): + """真空日(0 跌停)空文件无列——门面仍返标准列空表(schema 稳定)。""" + df = provider.get_limit_pool(kind="dtgc", date="2024-06-18") + assert df.empty + assert set(df.columns) == DTGC_KEYS + + def test_dtgc_real_columns_mapped(self, provider, tmp_path): + """dtgc 真实列(实测: 封单资金/连续跌停/开板次数...)→ 标准键。""" + f = tmp_path / "static" / "zt_pool_dtgc" / "20240620_zt_pool_dtgc.parquet" + pd.DataFrame([{ + "序号": 1, "代码": "600999", "名称": "跌停股", "涨跌幅": -10.0, + "最新价": 5.0, "成交额": 3e7, "流通市值": 1e9, "总市值": 2e9, + "动态市盈率": 15.0, "换手率": 8.0, "封单资金": 6e7, + "最后封板时间": "145950", "板上成交额": 1e6, "连续跌停": 2, + "开板次数": 1, "所属行业": "地产"}]).to_parquet(f, index=False) + df = provider.get_limit_pool(kind="dtgc", date="2024-06-20") + assert len(df) == 1 + assert set(df.columns) == DTGC_KEYS + row = df.iloc[0] + assert row["seal_amount"] == 6e7 + assert row["consecutive_limit_downs"] == 2 + assert row["open_count"] == 1 + assert row["pe_ttm"] == 15.0 + + def test_missing_code_column_fail_fast(self, provider, tmp_path): + """脏数据: 有行但缺 代码 → DataSchemaError(门面吸收列漂移的点)。""" + f = tmp_path / "static" / "zt_pool" / "20240620_zt_pool.parquet" + pd.DataFrame([{"名称": "无名股", "连板数": 1}]).to_parquet(f, index=False) + with pytest.raises(DataSchemaError, match="代码"): + provider.get_limit_pool(kind="zt", date="2024-06-20") + + def test_range_trading_filter_applies(self, provider): + df = provider.get_limit_pool(kind="zt", start="2024-06-18", + end="2024-06-20") + assert len(df) == 2 + assert "2024-06-19" not in set(df["trade_date"]) + + +# ======================== 注册表 ======================== + +class TestRegistry: + + def test_fetchers_registered(self): + assert FETCHERS["event_panel"] is EventPanelFetcher + assert FETCHERS["limit_pool"] is LimitPoolFetcher