- cross_sectional_zscore / composite_score / build_factor_panels 从 local_engine 迁入 quant/composite.py;新增统一入口 build_score_panel(daily, factor_specs) - local_engine re-export 保持旧引用兼容;selection/engine 的评分面板构建均指向 composite —— 选股与回测的复合分实现收敛于一处 - 回归:quant/eval/research/selection 一致性/qlib 引擎测试全过;全量 pytest 通过
293 lines
11 KiB
Python
293 lines
11 KiB
Python
"""LocalEngine —— 默认研究引擎(纯 pandas,AGENT.md §40 简单可替换优先)。
|
||
|
||
无未来函数纪律:
|
||
- 调仓日 t 的选股只使用 <=t 的因子值与收盘价
|
||
- 成交发生在 t 收盘(价格 = close[t] ± 滑点);t 当日组合收益用 t-1 收盘持仓结算,
|
||
调仓在 t 收盘生效、自 t+1 起计收益 —— 不存在「当日买入当日计收益」的未来函数
|
||
- 涨跌停 / 停牌约束按可达信息近似建模,未建模部分显式写入结果 unimplemented
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import math
|
||
from dataclasses import dataclass
|
||
from datetime import date
|
||
|
||
import pandas as pd
|
||
|
||
from app.domain.entities.research import (
|
||
BacktestResult,
|
||
BacktestSummary,
|
||
CurvePoint,
|
||
FactorTestReport,
|
||
MonthlyReturn,
|
||
Position,
|
||
ResearchSpec,
|
||
Trade,
|
||
YearlyReturn,
|
||
)
|
||
from app.quant.composite import ( # noqa: F401 —— re-export(模块化后旧引用仍可用)
|
||
build_factor_panels,
|
||
composite_score,
|
||
cross_sectional_zscore,
|
||
)
|
||
from app.quant.evaluation import run_factor_test
|
||
|
||
TRADING_DAYS = 252
|
||
_DEFAULT_UNIMPLEMENTED = [
|
||
"涨跌停按收盘价相对上一有效收盘近似判定(未建模开盘一字 / 集合竞价路径)",
|
||
"成交假设发生在调仓日收盘(未建模盘中价格路径与流动性冲击)",
|
||
]
|
||
|
||
|
||
def _limit_up_ratio(symbol: str) -> float:
|
||
"""按板块近似涨跌停幅度。"""
|
||
code = symbol[:3]
|
||
if code in {"300", "301", "688"}:
|
||
return 1.199
|
||
if code.startswith(("8", "4", "92")):
|
||
return 1.299
|
||
return 1.099
|
||
|
||
|
||
def rebalance_dates(index: pd.Index, rebalance: str, start: date) -> list[pd.Timestamp]:
|
||
"""按频率取首个交易日(>= start)。"""
|
||
periods = index.to_period("M" if rebalance == "monthly" else "W")
|
||
seen: dict = {}
|
||
order: list[pd.Timestamp] = []
|
||
for ts, per in zip(index, periods, strict=True):
|
||
if per not in seen:
|
||
seen[per] = ts
|
||
order.append(ts)
|
||
return [ts for ts in order if ts.date() >= start]
|
||
|
||
|
||
@dataclass
|
||
class EngineResult:
|
||
equity: pd.Series # index=date -> equity
|
||
trades: list[Trade]
|
||
positions: list[Position]
|
||
rebalance_notional: list[float]
|
||
|
||
|
||
class TopKBacktestRunner:
|
||
"""TopK 等权、固定调仓频率的低频回测。"""
|
||
|
||
def __init__(self, spec: ResearchSpec, score: pd.DataFrame, close: pd.DataFrame) -> None:
|
||
self.spec = spec
|
||
close = close.copy()
|
||
close.index = pd.to_datetime(close.index)
|
||
self.close = close.sort_index()
|
||
self.score = score.reindex(self.close.index).sort_index()
|
||
self.costs = spec.costs
|
||
# 上一有效收盘(用于涨跌停与收益结算,处理停牌日)
|
||
self.prev_close = self.close.ffill().shift(1)
|
||
|
||
def run(self) -> BacktestResult:
|
||
end_date = self.spec.period[1]
|
||
dates = [d for d in self.close.index if self.spec.period[0] <= d.date() <= end_date]
|
||
rebal = {
|
||
d
|
||
for d in rebalance_dates(self.close.index, self.spec.rebalance, self.spec.period[0])
|
||
if d.date() <= end_date
|
||
}
|
||
cash = float(self.spec.initial_capital)
|
||
shares: dict[str, float] = {}
|
||
entry_date: dict[str, date] = {}
|
||
entry_price: dict[str, float] = {}
|
||
equity_rows: dict[pd.Timestamp, float] = {}
|
||
trades: list[Trade] = []
|
||
positions: list[Position] = []
|
||
notional: list[float] = []
|
||
|
||
def _value(d: pd.Timestamp) -> float:
|
||
total = cash
|
||
for s, qty in shares.items():
|
||
if qty <= 0:
|
||
continue
|
||
px = self.close.at[d, s] if d in self.close.index else None
|
||
if px is None or (isinstance(px, float) and math.isnan(px)):
|
||
continue # 无行情日不计该仓(停牌近似,见 unimplemented)
|
||
total += float(qty * px)
|
||
return total
|
||
|
||
for d in dates:
|
||
if d in rebal:
|
||
cash = self._rebalance(
|
||
d, cash, shares, entry_date, entry_price, trades, positions, notional
|
||
)
|
||
equity_rows[d] = _value(d)
|
||
|
||
equity = pd.Series(equity_rows).sort_index()
|
||
return self._to_result(equity, trades, positions, notional)
|
||
|
||
# ---- 调仓(t 收盘执行,自 t+1 生效) ----
|
||
|
||
def _rebalance(self, d, cash, shares, entry_date, entry_price, trades, positions, notional):
|
||
close_d = self.close.loc[d]
|
||
prev_d = self.prev_close.loc[d]
|
||
sold_notional = 0.0
|
||
|
||
# 1) 卖出:跌停或无价(停牌)持仓保留,其余卖出
|
||
for s in [s for s in shares if shares[s] > 0]:
|
||
c, p = close_d[s], prev_d[s]
|
||
if _nan(c):
|
||
continue # 停牌无价:保留
|
||
if not _nan(p) and p > 0 and c / p <= 1.0 - (_limit_up_ratio(s) - 1.0):
|
||
continue # 跌停无法卖出:保留到下一调仓
|
||
qty = shares[s]
|
||
proceeds = qty * float(c) * (1 - self.costs.slippage_rate)
|
||
fee = proceeds * (self.costs.commission_rate + self.costs.stamp_tax_rate)
|
||
cash += proceeds - fee
|
||
sold_notional += proceeds
|
||
trades.append(
|
||
Trade(
|
||
entry_date=entry_date[s],
|
||
exit_date=d.date(),
|
||
symbol=s,
|
||
entry_price=entry_price[s],
|
||
exit_price=float(c),
|
||
return_pct=(float(c) / entry_price[s] - 1.0) * 100,
|
||
)
|
||
)
|
||
shares[s] = 0.0
|
||
entry_date.pop(s, None)
|
||
entry_price.pop(s, None)
|
||
|
||
# 2) 买入:取得分最高且可买的 TopN(涨停 / 无价剔除)
|
||
score_d = self.score.loc[d].dropna()
|
||
top = score_d.sort_values(ascending=False).index.tolist()
|
||
targets: list[str] = []
|
||
for s in top:
|
||
if len(targets) >= self.spec.selection.top_n:
|
||
break
|
||
c, p = close_d[s], prev_d[s]
|
||
if _nan(c) or _nan(p) or p <= 0:
|
||
continue
|
||
if c / p >= _limit_up_ratio(s):
|
||
continue # 涨停不可追买
|
||
targets.append(s)
|
||
|
||
if targets:
|
||
budget = cash / len(targets)
|
||
for s in targets:
|
||
c = float(close_d[s])
|
||
price_in = c * (1 + self.costs.slippage_rate)
|
||
invest = budget * (1 - self.costs.commission_rate)
|
||
shares[s] = invest / price_in
|
||
entry_date[s] = d.date()
|
||
entry_price[s] = price_in
|
||
notional.append(budget)
|
||
cash -= budget * len(targets)
|
||
|
||
# 3) 记录调仓后仓位
|
||
total = cash + sum(
|
||
float(self.close.at[d, s] * qty)
|
||
for s, qty in shares.items()
|
||
if qty > 0 and not _nan(self.close.at[d, s])
|
||
)
|
||
if total > 0:
|
||
for s, qty in shares.items():
|
||
if qty > 0 and not _nan(self.close.at[d, s]):
|
||
positions.append(
|
||
Position(
|
||
date=d.date(), symbol=s, weight=float(qty * self.close.at[d, s] / total)
|
||
)
|
||
)
|
||
return cash
|
||
|
||
# ---- 指标 ----
|
||
|
||
def _to_result(self, equity, trades, positions, notional) -> BacktestResult:
|
||
start, end = equity.index[0].date(), equity.index[-1].date()
|
||
init = float(self.spec.initial_capital)
|
||
final = float(equity.iloc[-1])
|
||
rets = equity.pct_change().dropna()
|
||
n = len(rets)
|
||
total_ret = (final / init - 1.0) * 100 if init else 0.0
|
||
annual = (
|
||
((final / init) ** (TRADING_DAYS / max(n, 1)) - 1.0) * 100
|
||
if final > 0 and init > 0
|
||
else -100.0
|
||
)
|
||
mean_r, std_r = (float(rets.mean()), float(rets.std(ddof=1))) if n else (0.0, 0.0)
|
||
sharpe = mean_r / std_r * math.sqrt(TRADING_DAYS) if std_r and mean_r else 0.0
|
||
vol = std_r * math.sqrt(TRADING_DAYS) * 100
|
||
dd = (equity / equity.cummax() - 1.0).min() * 100
|
||
wins = [t for t in trades if t.return_pct > 0]
|
||
win_rate = len(wins) / len(trades) * 100 if trades else 0.0
|
||
avg_turn = (sum(notional) / len(notional) / ((init + final) / 2)) * 100 if notional else 0.0
|
||
|
||
eq_pts = [CurvePoint(date=d.date(), value=round(float(v), 2)) for d, v in equity.items()]
|
||
dd_series = (equity / equity.cummax() - 1.0) * 100
|
||
drawdown = [
|
||
CurvePoint(date=d.date(), value=round(float(v), 3)) for d, v in dd_series.items()
|
||
]
|
||
|
||
monthly: list[MonthlyReturn] = []
|
||
yearly: list[YearlyReturn] = []
|
||
if len(equity) > 1:
|
||
m = equity.resample("ME").last().pct_change().dropna()
|
||
monthly = [
|
||
MonthlyReturn(
|
||
year=int(d.year), month=int(d.month), return_pct=round(float(v) * 100, 3)
|
||
)
|
||
for d, v in m.items()
|
||
]
|
||
y = equity.resample("YE").last().pct_change().dropna()
|
||
yearly = [
|
||
YearlyReturn(year=int(d.year), return_pct=round(float(v) * 100, 3))
|
||
for d, v in y.items()
|
||
]
|
||
|
||
summary = BacktestSummary(
|
||
start=start,
|
||
end=end,
|
||
initial_capital=round(init, 2),
|
||
final_equity=round(final, 2),
|
||
total_return_pct=round(total_ret, 3),
|
||
annual_return_pct=round(annual, 3),
|
||
sharpe=round(sharpe, 3),
|
||
max_drawdown_pct=round(float(dd), 3),
|
||
volatility_pct=round(vol, 3),
|
||
win_rate_pct=round(win_rate, 2),
|
||
total_trades=len(trades),
|
||
avg_turnover_pct=round(avg_turn, 2),
|
||
)
|
||
return BacktestResult(
|
||
summary=summary,
|
||
equity_curve=eq_pts,
|
||
drawdown=drawdown,
|
||
monthly_returns=monthly,
|
||
yearly_returns=yearly,
|
||
positions=positions,
|
||
trades=trades,
|
||
turnover_pct=round(sum(notional) / max(init, 1) * 100, 2),
|
||
unimplemented=list(_DEFAULT_UNIMPLEMENTED),
|
||
config_snapshot=self.spec.model_dump(mode="json"),
|
||
)
|
||
|
||
|
||
def run_spec_factor_test(
|
||
daily: pd.DataFrame,
|
||
spec: ResearchSpec,
|
||
horizon_days: int = 21,
|
||
) -> tuple[FactorTestReport, dict[str, pd.DataFrame]]:
|
||
"""单因子测试:因子面板 + 未来 horizon 收益 → FactorTestReport。"""
|
||
assert spec.type == "factor_test"
|
||
factor_name = spec.factors[0].name
|
||
panels = build_factor_panels(daily, spec.factors)
|
||
panel = panels[0][1]
|
||
close = daily.pivot(index="trade_date", columns="symbol", values="close").sort_index()
|
||
close.index = pd.to_datetime(close.index)
|
||
forward = close.shift(-horizon_days) / close - 1.0
|
||
report = run_factor_test(panel, forward, factor_name=factor_name)
|
||
return report, {factor_name: panel}
|
||
|
||
|
||
def _nan(v) -> bool:
|
||
try:
|
||
return bool(math.isnan(float(v)))
|
||
except (TypeError, ValueError):
|
||
return False
|