feat(provider): A档截面数据使用层出口——泛型底座+涨停池门面双层接口
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:
@@ -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",
|
||||
]
|
||||
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user