feat(provider): A档截面数据使用层出口——泛型底座+涨停池门面双层接口
CI/CD / test (push) Failing after 3s
CI/CD / nas-deploy (push) Has been skipped
CI/CD / nas-verify (push) Has been skipped

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]
This commit is contained in:
2026-09-02 10:43:11 +08:00
parent 8c6a3f90d0
commit 4966d30a11
7 changed files with 669 additions and 2 deletions
@@ -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
@@ -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",
]
+70 -1
View File
@@ -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,
*,
@@ -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)
@@ -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)
@@ -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]],
@@ -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