feat(minute): 分钟红7新增「N日内涨停过」条件 + minute_filter 日线窗口契约

分钟策略此前只能访问当日分钟窗口, 无法叠加日线维度条件。本次为
minute_filter 后端扩展可选日线历史契约:

- 契约: 策略声明 META["daily_history_bars"] (0-250) + filter_minute_history
  接受 daily 关键字 (加载期校验, 纯分钟策略零改动); 引擎聚合各策略声明
  (minute_daily_history_bars) 后由 ScreenerService 1m 分支装配
  context.daily_history (enriched 日线窗口), run() 以 daily= 注入
- 分钟红7: require_limit_up (默认开) + limit_up_days (5-60, 默认20) 参数;
  涨停判定复用 enriched 预计算信号 signal_limit_up (收盘封板) 或
  signal_broken_limit_up (炸板盘中触及), 任一命中即算涨停过, 输出
  recent_limit_ups 次数列; 日线窗口缺失时失败闭合 (宁可漏过不可错报)
- 测试 25 项: 涨停信号过滤/炸板计数/回看窗口边界(第20日含第21日不含)/
  失败闭合/开关旁路/加载校验(缺 daily 关键字与超范围)/引擎注入/服务装配

实盘验证: 用户参数(bars=6)基线 24 只 → 要求涨停过后 6 只, 全部带
recent_limit_ups; 全量 1112 项后端测试通过。
This commit is contained in:
shy3130
2026-08-30 19:05:09 +08:00
parent bfa87d1f0b
commit b82c4eaedb
4 changed files with 285 additions and 15 deletions
+8
View File
@@ -390,12 +390,20 @@ class ScreenerService:
# 分钟策略数据源是本地当日分钟K分区 (单分区文件直读), 与日线
# enriched 历史窗口无关, 不走 required_history_bars 日线路径。
history = self._load_minute_history(as_of, current)
# 策略声明 META["daily_history_bars"] 时额外装配日线 enriched 窗口,
# 供分钟策略叠加日线维度条件 (如 N 日内涨停过)。
daily_history = None
if engine is not None:
daily_bars = engine.minute_daily_history_bars(strategy_ids)
if daily_bars > 0:
daily_history = self._load_enriched_history(as_of, daily_bars)
return StrategyDataContext(
asset_type=self.asset_type,
timeframe=timeframe,
as_of=as_of,
current=current,
history=history,
daily_history=daily_history,
market=None,
cache_key=cache_key,
)
@@ -4,6 +4,11 @@
(symbol, datetime, open, high, low, close, volume, amount),
由 ScreenerService.build_strategy_context 的 1m 分支从本地 kline_minute
分区注入; 策略本身不感知数据来源 (本地同步 / 盘中增量刷新对它透明)。
META["daily_history_bars"] 声明叠加日线维度的条件 (N 日内涨停过):
引擎会以 daily= 关键字注入日线 enriched 窗口, 涨停判定直接复用
enriched 预计算信号 — signal_limit_up (收盘封板) 或 signal_broken_limit_up
(炸板: 盘中触及涨停未封住), 任一命中即算"盘中涨停过"
"""
import polars as pl
@@ -11,10 +16,12 @@ import polars as pl
META = {
"id": "minute_red_streak",
"name": "分钟红7",
"description": "开盘前7根1分钟K至少5根收红, 最高的2根(按最高价)都是红K",
"description": "开盘前7根1分钟K至少5根收红, 最高的2根(按最高价)都是红K, 且近20日盘中触及过涨停",
"tags": ["分钟", "形态", "短线"],
"asset_types": ["stock"],
"timeframes": ["1m"],
# 日线 enriched 窗口 (交易日语义, 含 as_of): 覆盖 limit_up_days 参数上限
"daily_history_bars": 60,
"params": [
{
"id": "bars",
@@ -49,6 +56,21 @@ META = {
"type": "bool",
"default": False,
},
{
"id": "require_limit_up",
"label": "要求N日内涨停过",
"type": "bool",
"default": True,
},
{
"id": "limit_up_days",
"label": "涨停回看天数",
"type": "int",
"default": 20,
"min": 5,
"max": 60,
"step": 1,
},
],
"order_by": "red_count",
"descending": True,
@@ -60,12 +82,41 @@ ENTRY_SIGNALS: list[str] = []
EXIT_SIGNALS: list[str] = []
def filter_minute_history(df: pl.DataFrame, params: dict) -> pl.DataFrame:
def _recent_limit_ups(daily: pl.DataFrame | None, lookback: int) -> pl.DataFrame:
"""日线窗口 → (symbol, recent_limit_ups) 近 lookback 个交易日的涨停次数。
涨停过 = signal_limit_up (收盘封板) 或 signal_broken_limit_up (炸板触及)。
日线窗口缺失 / 无涨停信号列 → 返回空表 (调用方 inner join 即失败闭合,
宁可漏过不可错报)。
"""
empty = pl.DataFrame(schema={"symbol": pl.Utf8, "recent_limit_ups": pl.UInt32})
if daily is None or daily.is_empty():
return empty
if not {"signal_limit_up", "signal_broken_limit_up"}.issubset(daily.columns):
return empty
return (
daily.select("symbol", "date", "signal_limit_up", "signal_broken_limit_up")
.sort(["symbol", "date"])
.filter(pl.int_range(pl.len()).over("symbol") >= pl.len().over("symbol") - lookback)
.group_by("symbol")
.agg(
recent_limit_ups=(
pl.col("signal_limit_up").fill_null(False)
| pl.col("signal_broken_limit_up").fill_null(False)
).sum()
)
.filter(pl.col("recent_limit_ups") > 0)
)
def filter_minute_history(df: pl.DataFrame, params: dict, *, daily: pl.DataFrame | None = None) -> pl.DataFrame:
"""红K形态过滤: 全向量化, 无逐行 Python 循环。
- 每标的按时间取当日最早 bars 根 (开盘窗口); 不足 bars 根不触发
- 红 = close > open; 窗口内红K数 >= min_red
- 按 rank_by (high / close) 降序取前 top_red 根, 同值取时间更晚者, 需全红
- require_limit_up: 近 limit_up_days 个交易日盘中触及过涨停 (日线维度,
由 daily 窗口的预计算涨停信号判定; 窗口缺失时失败闭合不触发)
"""
bars = int(params.get("bars") or 7)
min_red = min(int(params.get("min_red") or 5), bars)
@@ -99,7 +150,7 @@ def filter_minute_history(df: pl.DataFrame, params: dict) -> pl.DataFrame:
.agg(top_red_count=pl.col("_red").sum())
)
return (
result = (
window.join(top, on="symbol", how="inner")
.filter(
(pl.col("bars_checked") >= bars)
@@ -108,3 +159,7 @@ def filter_minute_history(df: pl.DataFrame, params: dict) -> pl.DataFrame:
)
.drop("bars_checked")
)
if params.get("require_limit_up", True):
lookback = max(5, min(int(params.get("limit_up_days") or 20), 60))
result = result.join(_recent_limit_ups(daily, lookback), on="symbol", how="inner")
return result
+38 -1
View File
@@ -155,6 +155,9 @@ class StrategyDataContext:
as_of: date
current: pl.DataFrame | None = None
history: pl.DataFrame | None = None
# 仅 1m 分支: 策略声明 META["daily_history_bars"] 时注入的日线 enriched 窗口,
# 供分钟策略叠加日线维度条件 (如 N 日内涨停过); 未声明时为 None。
daily_history: pl.DataFrame | None = None
market: Any | None = None
cache_key: str | None = None
@@ -201,6 +204,9 @@ class StrategyDef:
composite: CompositeSpec | None = None # 仅 backend=="composite" 时非空
# 仅 backend=="minute_filter" 时非空: 输入为当日分钟K窗口, 输出为命中标的行
filter_minute_history_fn: Callable[[pl.DataFrame, dict], pl.DataFrame] | None = None
# 仅 minute_filter: META["daily_history_bars"] 声明需要的日线历史窗口 (0=不需要;
# >0 时 filter_minute_history 必须接受 daily 关键字, 引擎注入 context.daily_history)
minute_daily_bars: int = 0
@dataclass
@@ -493,6 +499,7 @@ class StrategyEngine:
matrix_strategy = getattr(mod, "MATRIX_STRATEGY", None)
composite_spec: CompositeSpec | None = None
minute_daily_bars = 0
if execution_backend == "matrix_native":
from app.backtest.matrix import MatrixStrategy
@@ -535,6 +542,22 @@ class StrategyEngine:
raise ValueError(
"minute_filter strategy must declare timeframes == ['1m']"
)
# 可选日线历史窗口: 声明 daily_history_bars 时 fn 必须接受 daily 关键字,
# 引擎会把 context.daily_history (enriched 日线窗口) 注入进来。
minute_daily_bars = int(meta.get("daily_history_bars") or 0)
if minute_daily_bars < 0 or minute_daily_bars > 250:
raise ValueError(
"minute_filter daily_history_bars must be within [0, 250]"
)
if minute_daily_bars > 0:
import inspect
sig = inspect.signature(filter_minute_history_fn)
if "daily" not in sig.parameters:
raise ValueError(
"minute_filter daily_history_bars requires "
"filter_minute_history to accept a 'daily' keyword"
)
elif filter_history_fn is None or filter_fn is not None:
raise ValueError("python_history_legacy strategy must declare only filter_history")
@@ -559,6 +582,7 @@ class StrategyEngine:
matrix_strategy=matrix_strategy,
composite=composite_spec,
filter_minute_history_fn=filter_minute_history_fn,
minute_daily_bars=minute_daily_bars,
)
def reload(self) -> None:
@@ -673,6 +697,16 @@ class StrategyEngine:
return None
return max(0, int(value))
def minute_daily_history_bars(self, strategy_ids: list[str]) -> int:
"""1m 分支需要的日线 enriched 窗口大小: 各 minute_filter 策略声明的
META["daily_history_bars"] 取 max, 未声明 (纯分钟策略) 为 0。"""
required = 0
for strategy_id in strategy_ids:
strategy = self.get(strategy_id)
if strategy.execution_backend == "minute_filter":
required = max(required, strategy.minute_daily_bars)
return required
def required_history_bars(
self,
strategy_ids: list[str],
@@ -917,7 +951,10 @@ class StrategyEngine:
strategy_id=strategy_id,
exit_signal_hits=exit_signal_hits,
)
df = s.filter_minute_history_fn(history, params)
if s.minute_daily_bars > 0:
df = s.filter_minute_history_fn(history, params, daily=context.daily_history)
else:
df = s.filter_minute_history_fn(history, params)
# 基础过滤/展示列 (name/total_shares/change_pct 等) 来自 enriched 快照,
# 在命中结果上事后联表, 避免把 enriched 列铺到全市场分钟行上。
if current is not None and not current.is_empty():