feat: add post-execution daily ledger

This commit is contained in:
ao gong
2026-08-21 22:01:35 +08:00
parent a44e2d306a
commit b794ab2e8f
2 changed files with 414 additions and 109 deletions
+293 -109
View File
@@ -22,7 +22,7 @@ from __future__ import annotations
import math import math
from collections.abc import Mapping from collections.abc import Mapping
from dataclasses import dataclass from dataclasses import dataclass, replace
from typing import Any from typing import Any
import pandas as pd import pandas as pd
@@ -103,6 +103,9 @@ class ExecutionResult:
net_cash_flow: float # 净现金流(买入为负,卖出为正) net_cash_flow: float # 净现金流(买入为负,卖出为正)
partial_fill_pct: float = 1.0 # 实际成交占目标的比例(1.0 = 全部成交) partial_fill_pct: float = 1.0 # 实际成交占目标的比例(1.0 = 全部成交)
blocked_reason: str = "" # 阻塞原因(如涨跌停停牌) blocked_reason: str = "" # 阻塞原因(如涨跌停停牌)
side: str = "" # buy / sell;未成交记录也保留目标方向
quantity: float = 0.0 # 实际成交股数
price: float = 0.0 # 未含滑点的参考执行价
def _apply_costs( def _apply_costs(
@@ -369,6 +372,56 @@ class ExecutionSimulationResult:
dtype=float, dtype=float,
) )
@property
def normalized_nav_series(self) -> pd.Series:
"""返回以初始资金为 1 的净值曲线副本。"""
nav = self.nav_series
if self.initial_cash == 0:
return pd.Series(0.0, index=nav.index, dtype=float)
return nav / self.initial_cash
@property
def daily_returns(self) -> pd.Series:
"""返回逐日收益;首日相对初始资金计算,保留首日交易成本。"""
nav = self.nav_series
if nav.empty:
return nav
returns = nav.pct_change()
returns.iloc[0] = (
nav.iloc[0] / self.initial_cash - 1.0 if self.initial_cash != 0 else 0.0
)
return returns.fillna(0.0)
@property
def trades_frame(self) -> pd.DataFrame:
"""返回可投影到平台成交明细的实际成交表,不包含纯拒绝记录。"""
columns = [
"trade_date",
"ts_code",
"side",
"qty",
"price",
"amount",
"fee",
"slippage",
]
rows = [
{
"trade_date": daily.date,
"ts_code": execution.stock_code,
"side": execution.side,
"qty": execution.quantity,
"price": execution.price,
"amount": execution.executed_value,
"fee": execution.commission + execution.stamp_tax,
"slippage": execution.slippage_cost,
}
for daily in self.daily_executions
for execution in daily.executions
if execution.quantity > 0
]
return pd.DataFrame(rows, columns=columns)
@property @property
def total_costs(self) -> float: def total_costs(self) -> float:
"""汇总实际成交产生的成本。""" """汇总实际成交产生的成本。"""
@@ -420,6 +473,7 @@ def _blocked_execution(stock_code: str, target_value: float, reason: str) -> Exe
net_cash_flow=0.0, net_cash_flow=0.0,
partial_fill_pct=0.0, partial_fill_pct=0.0,
blocked_reason=reason, blocked_reason=reason,
side="buy" if target_value > 0 else "sell" if target_value < 0 else "",
) )
@@ -466,6 +520,236 @@ def _partially_fill_buy(
) )
def _rebalance_at_prices(
date: str,
targets: Mapping[str, float],
prices: Mapping[str, float],
cash: float,
holdings: dict[str, float],
config: ExecutionConfig,
) -> tuple[float, tuple[ExecutionResult, ...], float, float]:
"""在单一执行时点按目标权重差额调仓,并原地更新 holdings。"""
normalized_targets = _validate_target_weights(date, targets)
for held_code in holdings:
held_price = prices.get(held_code)
if held_price is None or not math.isfinite(held_price) or held_price <= 0:
raise ValueError(f"missing price for held asset {held_code} on {date!r}")
nav_before = cash + sum(
shares * prices.get(stock_code, 0.0)
for stock_code, shares in holdings.items()
)
effective_targets = dict.fromkeys(holdings, 0.0)
effective_targets.update(normalized_targets)
buy_weights: dict[str, float] = {}
sell_weights: dict[str, float] = {}
rejected: list[ExecutionResult] = []
for stock_code, target_weight in effective_targets.items():
price = prices.get(stock_code)
target_value = float(target_weight) * nav_before
if price is None or not math.isfinite(price) or price <= 0:
if target_value != 0 or holdings.get(stock_code, 0.0) != 0:
rejected.append(_blocked_execution(stock_code, target_value, "missing_price"))
continue
current_value = holdings.get(stock_code, 0.0) * price
trade_value = target_value - current_value
if abs(trade_value) < config.min_trade_amount or math.isclose(
trade_value, 0.0, abs_tol=1e-12
):
continue
if nav_before == 0:
rejected.append(_blocked_execution(stock_code, trade_value, "zero_nav"))
continue
destination = buy_weights if trade_value > 0 else sell_weights
destination[stock_code] = trade_value / nav_before
sell_executions = simulate_execution(sell_weights, nav_before, config)
filled: list[ExecutionResult] = []
for raw_execution in sell_executions:
price = prices[raw_execution.stock_code]
quantity = abs(raw_execution.target_value) / price
execution = replace(
raw_execution,
side="sell",
quantity=quantity,
price=price,
)
held = holdings.get(execution.stock_code, 0.0)
holdings[execution.stock_code] = max(0.0, held - quantity)
if holdings[execution.stock_code] < 1e-6:
del holdings[execution.stock_code]
cash += execution.net_cash_flow
filled.append(execution)
desired_buys = simulate_execution(buy_weights, nav_before, config)
required_cash = sum(-execution.net_cash_flow for execution in desired_buys)
buy_fill_pct = min(1.0, max(cash, 0.0) / required_cash) if required_cash > 0 else 1.0
for desired in desired_buys:
if buy_fill_pct == 0:
rejected.append(
_blocked_execution(desired.stock_code, desired.target_value, "insufficient_cash")
)
continue
raw_execution = (
desired
if buy_fill_pct == 1.0
else _partially_fill_buy(desired, buy_fill_pct, config)
)
price = prices[raw_execution.stock_code]
quantity = abs(raw_execution.target_value) * raw_execution.partial_fill_pct / price
execution = replace(
raw_execution,
side="buy",
quantity=quantity,
price=price,
)
holdings[execution.stock_code] = holdings.get(execution.stock_code, 0.0) + quantity
cash += execution.net_cash_flow
if math.isclose(cash, 0.0, abs_tol=1e-9):
cash = 0.0
filled.append(execution)
executions = (*filled, *rejected)
nav_after = cash + sum(
shares * prices.get(stock_code, 0.0)
for stock_code, shares in holdings.items()
)
return cash, executions, nav_before, nav_after
def _validate_sparse_daily_histories(
target_weights_history: list[tuple[str, dict[str, float]]],
execution_price_history: list[tuple[str, dict[str, float]]],
valuation_price_history: list[tuple[str, dict[str, float]]],
) -> tuple[
dict[str, dict[str, float]],
dict[str, dict[str, float]],
list[tuple[str, dict[str, float]]],
]:
"""校验稀疏调仓与完整估值日历,并隔离调用方可变输入。"""
target_dates = [date for date, _ in target_weights_history]
execution_dates = [date for date, _ in execution_price_history]
valuation_dates = [date for date, _ in valuation_price_history]
if len(set(target_dates)) != len(target_dates):
raise ValueError("target_weights_history must contain unique dates")
if len(set(execution_dates)) != len(execution_dates):
raise ValueError("execution_price_history must contain unique dates")
if len(set(valuation_dates)) != len(valuation_dates):
raise ValueError("valuation_price_history must contain unique dates")
if execution_dates != target_dates:
raise ValueError("execution price dates must exactly match target weight dates")
valuation_positions = {date: index for index, date in enumerate(valuation_dates)}
missing_dates = [date for date in target_dates if date not in valuation_positions]
if missing_dates:
raise ValueError(f"target dates must belong to valuation calendar: {missing_dates}")
positions = [valuation_positions[date] for date in target_dates]
if positions != sorted(positions):
raise ValueError("target weights must follow valuation calendar order")
targets = {date: dict(values) for date, values in target_weights_history}
execution_prices = {date: dict(values) for date, values in execution_price_history}
valuation_prices = [(date, dict(values)) for date, values in valuation_price_history]
return targets, execution_prices, valuation_prices
def _simulate_daily_ledger(
target_weights_history: list[tuple[str, dict[str, float]]],
execution_price_history: list[tuple[str, dict[str, float]]],
valuation_price_history: list[tuple[str, dict[str, float]]],
initial_cash: float,
config: ExecutionConfig,
) -> ExecutionSimulationResult:
targets_by_date, execution_prices_by_date, valuation_history = (
_validate_sparse_daily_histories(
target_weights_history,
execution_price_history,
valuation_price_history,
)
)
cash = initial_cash
holdings: dict[str, float] = {}
positions: list[DailyPosition] = []
daily_executions: list[DailyExecution] = []
for date, valuation_prices in valuation_history:
targets = targets_by_date.get(date)
if targets is None:
executions: tuple[ExecutionResult, ...] = ()
nav_before = 0.0
nav_after = 0.0
rebalance_triggered = False
else:
cash, executions, nav_before, nav_after = _rebalance_at_prices(
date,
targets,
execution_prices_by_date[date],
cash,
holdings,
config,
)
rebalance_triggered = any(execution.quantity > 0 for execution in executions)
for held_code in holdings:
valuation_price = valuation_prices.get(held_code)
if (
valuation_price is None
or not math.isfinite(valuation_price)
or valuation_price <= 0
):
raise ValueError(
f"missing valuation price for held asset {held_code} on {date!r}"
)
portfolio_value = cash + sum(
shares * valuation_prices[stock_code]
for stock_code, shares in holdings.items()
)
if targets is None:
nav_before = portfolio_value
nav_after = portfolio_value
positions.append(DailyPosition(date, cash, dict(holdings), portfolio_value))
daily_executions.append(
DailyExecution(
date=date,
executions=executions,
nav_before=nav_before,
nav_after=nav_after,
rebalance_triggered=rebalance_triggered,
)
)
return ExecutionSimulationResult(
initial_cash=initial_cash,
positions=tuple(positions),
daily_executions=tuple(daily_executions),
)
def simulate_daily_ledger_with_audit(
target_weights_history: list[tuple[str, dict[str, float]]],
execution_price_history: list[tuple[str, dict[str, float]]],
valuation_price_history: list[tuple[str, dict[str, float]]],
initial_cash: float,
config: ExecutionConfig | None = None,
) -> ExecutionSimulationResult:
"""以稀疏调仓和完整日历运行成交后持仓 Ledger。
执行价只用于调仓日现金与股数变化,估值价用于每个交易日日末 NAV;二者
显式分离,从而支持“下一日 open 成交、同日 close 估值”的无前视研究。
"""
if not math.isfinite(initial_cash) or initial_cash <= 0:
raise ValueError(f"initial_cash must be positive and finite, got {initial_cash}")
return _simulate_daily_ledger(
target_weights_history,
execution_price_history,
valuation_price_history,
initial_cash,
ExecutionConfig() if config is None else config,
)
def simulate_multi_day_with_audit( def simulate_multi_day_with_audit(
target_weights_history: list[tuple[str, dict[str, float]]], target_weights_history: list[tuple[str, dict[str, float]]],
price_history: list[tuple[str, dict[str, float]]], price_history: list[tuple[str, dict[str, float]]],
@@ -488,122 +772,21 @@ def simulate_multi_day_with_audit(
- 每日先按当日 close 估值,再交易“目标市值 - 当前市值”的差额 - 每日先按当日 close 估值,再交易“目标市值 - 当前市值”的差额
- 此处简化为当日 close 成交;调用方必须传入已正确滞后的目标权重 - 此处简化为当日 close 成交;调用方必须传入已正确滞后的目标权重
""" """
if config is None:
config = ExecutionConfig()
if not math.isfinite(initial_cash) or initial_cash < 0: if not math.isfinite(initial_cash) or initial_cash < 0:
raise ValueError(f"initial_cash must be finite and non-negative, got {initial_cash}") raise ValueError(f"initial_cash must be finite and non-negative, got {initial_cash}")
if len(target_weights_history) != len(price_history): if len(target_weights_history) != len(price_history):
raise ValueError("target_weights_history and price_history must have same length") raise ValueError("target_weights_history and price_history must have same length")
if not target_weights_history: for (date, _), (price_date, _) in zip(target_weights_history, price_history, strict=True):
return ExecutionSimulationResult(initial_cash, (), ())
cash = initial_cash
holdings: dict[str, float] = {}
positions: list[DailyPosition] = []
daily_executions: list[DailyExecution] = []
for (date, targets), (price_date, prices) in zip(
target_weights_history, price_history, strict=True
):
if date != price_date: if date != price_date:
raise ValueError( raise ValueError(
f"target and price dates must match, got {date!r} and {price_date!r}" f"target and price dates must match, got {date!r} and {price_date!r}"
) )
return _simulate_daily_ledger(
normalized_targets = _validate_target_weights(date, targets) target_weights_history,
for held_code in holdings: price_history,
held_price = prices.get(held_code) price_history,
if held_price is None or not math.isfinite(held_price) or held_price <= 0: initial_cash,
raise ValueError(f"missing price for held asset {held_code} on {date!r}") ExecutionConfig() if config is None else config,
nav_before = cash + sum(
shares * prices.get(stock_code, 0.0)
for stock_code, shares in holdings.items()
)
effective_targets = dict.fromkeys(holdings, 0.0)
effective_targets.update(normalized_targets)
buy_weights: dict[str, float] = {}
sell_weights: dict[str, float] = {}
rejected: list[ExecutionResult] = []
for stock_code, target_weight in effective_targets.items():
price = prices.get(stock_code)
target_value = float(target_weight) * nav_before
if price is None or not math.isfinite(price) or price <= 0:
if target_value != 0 or holdings.get(stock_code, 0.0) != 0:
rejected.append(_blocked_execution(stock_code, target_value, "missing_price"))
continue
current_value = holdings.get(stock_code, 0.0) * price
trade_value = target_value - current_value
if abs(trade_value) < config.min_trade_amount or math.isclose(
trade_value, 0.0, abs_tol=1e-12
):
continue
if nav_before == 0:
rejected.append(_blocked_execution(stock_code, trade_value, "zero_nav"))
continue
destination = buy_weights if trade_value > 0 else sell_weights
destination[stock_code] = trade_value / nav_before
sell_executions = simulate_execution(sell_weights, nav_before, config)
filled: list[ExecutionResult] = []
for execution in sell_executions:
price = prices[execution.stock_code]
share_change = abs(execution.target_value) / price
held = holdings.get(execution.stock_code, 0.0)
holdings[execution.stock_code] = max(0.0, held - share_change)
if holdings[execution.stock_code] < 1e-6:
del holdings[execution.stock_code]
cash += execution.net_cash_flow
filled.append(execution)
desired_buys = simulate_execution(buy_weights, nav_before, config)
required_cash = sum(-execution.net_cash_flow for execution in desired_buys)
buy_fill_pct = min(1.0, max(cash, 0.0) / required_cash) if required_cash > 0 else 1.0
for desired in desired_buys:
if buy_fill_pct == 0:
rejected.append(
_blocked_execution(desired.stock_code, desired.target_value, "insufficient_cash")
)
continue
execution = (
desired
if buy_fill_pct == 1.0
else _partially_fill_buy(desired, buy_fill_pct, config)
)
price = prices[execution.stock_code]
share_change = (
execution.target_value * execution.partial_fill_pct / price
)
holdings[execution.stock_code] = (
holdings.get(execution.stock_code, 0.0) + share_change
)
cash += execution.net_cash_flow
if math.isclose(cash, 0.0, abs_tol=1e-9):
cash = 0.0
filled.append(execution)
executions = (*filled, *rejected)
nav_after = cash + sum(
shares * prices.get(stock_code, 0.0)
for stock_code, shares in holdings.items()
)
positions.append(DailyPosition(date, cash, dict(holdings), nav_after))
daily_executions.append(
DailyExecution(
date=date,
executions=executions,
nav_before=nav_before,
nav_after=nav_after,
rebalance_triggered=bool(filled),
)
)
return ExecutionSimulationResult(
initial_cash=initial_cash,
positions=tuple(positions),
daily_executions=tuple(daily_executions),
) )
@@ -785,6 +968,7 @@ __all__ = [
"DailyPosition", "DailyPosition",
"DailyExecution", "DailyExecution",
"ExecutionSimulationResult", "ExecutionSimulationResult",
"simulate_daily_ledger_with_audit",
"simulate_multi_day", "simulate_multi_day",
"simulate_multi_day_with_audit", "simulate_multi_day_with_audit",
"run_end_to_end_poc", "run_end_to_end_poc",
+121
View File
@@ -7,6 +7,7 @@
from __future__ import annotations from __future__ import annotations
from collections.abc import Mapping
from dataclasses import dataclass from dataclasses import dataclass
import numpy as np import numpy as np
@@ -16,15 +17,19 @@ from pandas.api.types import is_numeric_dtype
from quant_engine.execution import ( from quant_engine.execution import (
ExecutionConfig, ExecutionConfig,
ExecutionSimulationResult, ExecutionSimulationResult,
simulate_daily_ledger_with_audit,
simulate_multi_day_with_audit, simulate_multi_day_with_audit,
) )
from quant_engine.metrics import summary as metrics_summary
from quant_engine.portfolio_construction import scores_to_weight_table from quant_engine.portfolio_construction import scores_to_weight_table
__all__ = [ __all__ = [
"TargetWeightSchedule", "TargetWeightSchedule",
"FactorExecutionResult", "FactorExecutionResult",
"FactorBacktestResult",
"schedule_target_weights", "schedule_target_weights",
"run_factor_execution_research", "run_factor_execution_research",
"run_factor_backtest_research",
] ]
@@ -49,6 +54,41 @@ class FactorExecutionResult:
execution: ExecutionSimulationResult execution: ExecutionSimulationResult
@dataclass(frozen=True, slots=True, eq=False)
class FactorBacktestResult:
"""因子、成交后日频 Ledger 与绩效的一次可复现快照。"""
factor_scores: pd.DataFrame
execution_prices: pd.DataFrame
valuation_prices: pd.DataFrame
schedule: TargetWeightSchedule
execution_price_field: str
valuation_price_field: str
execution: ExecutionSimulationResult
@property
def nav(self) -> pd.Series:
"""返回以初始资金归一化为 1 的日频 NAV。"""
return pd.Series(
self.execution.normalized_nav_series.to_numpy(copy=True),
index=self.valuation_prices.index.copy(),
name="nav",
)
@property
def returns(self) -> pd.Series:
"""返回包含首日成本影响的日频收益。"""
return pd.Series(
self.execution.daily_returns.to_numpy(copy=True),
index=self.valuation_prices.index.copy(),
name="returns",
)
def stats(self, rf: float = 0.0) -> Mapping[str, float]:
"""复用标准绩效口径计算指标。"""
return metrics_summary(self.returns, rf)
def _validate_datetime_index(index: pd.Index, name: str) -> pd.DatetimeIndex: def _validate_datetime_index(index: pd.Index, name: str) -> pd.DatetimeIndex:
if not isinstance(index, pd.DatetimeIndex): if not isinstance(index, pd.DatetimeIndex):
raise TypeError(f"{name} must use a DatetimeIndex") raise TypeError(f"{name} must use a DatetimeIndex")
@@ -215,3 +255,84 @@ def run_factor_execution_research(
execution_price_field=price_field, execution_price_field=price_field,
execution=execution, execution=execution,
) )
def run_factor_backtest_research(
factor_scores: pd.DataFrame,
execution_prices: pd.DataFrame,
valuation_prices: pd.DataFrame,
*,
top_k: int,
execution_price_field: str,
valuation_price_field: str,
lag_sessions: int = 1,
gross_exposure: float = 1.0,
largest: bool = True,
initial_cash: float = 1_000_000.0,
config: ExecutionConfig | None = None,
) -> FactorBacktestResult:
"""运行 PIT 因子到成交后日频 Ledger、收益与绩效的可信研究链路。"""
execution_field = execution_price_field.strip()
valuation_field = valuation_price_field.strip()
if not execution_field:
raise ValueError("execution_price_field must be non-empty")
if not valuation_field:
raise ValueError("valuation_price_field must be non-empty")
execution_calendar = _validate_execution_prices(execution_prices)
valuation_calendar = _validate_execution_prices(valuation_prices)
if not execution_calendar.equals(valuation_calendar):
raise ValueError("execution and valuation prices must use matching trading calendars")
if not execution_prices.columns.equals(valuation_prices.columns):
raise ValueError("execution and valuation prices must use matching asset labels")
factor_snapshot = factor_scores.copy(deep=True)
execution_snapshot = execution_prices.copy(deep=True)
valuation_snapshot = valuation_prices.copy(deep=True)
decision_weights = scores_to_weight_table(
factor_snapshot,
top_k,
gross_exposure=gross_exposure,
largest=largest,
)
schedule = schedule_target_weights(
decision_weights,
execution_calendar,
lag_sessions=lag_sessions,
)
target_history: list[tuple[str, dict[str, float]]] = []
execution_history: list[tuple[str, dict[str, float]]] = []
for execution_date, weights in schedule.execution_weights.iterrows():
date_label = str(pd.Timestamp(execution_date))
target_history.append(
(date_label, {asset: float(weight) for asset, weight in weights.items()})
)
prices = execution_snapshot.loc[execution_date]
execution_history.append(
(date_label, {asset: float(price) for asset, price in prices.items()})
)
valuation_history = [
(
str(pd.Timestamp(valuation_date)),
{asset: float(price) for asset, price in prices.items()},
)
for valuation_date, prices in valuation_snapshot.iterrows()
]
execution = simulate_daily_ledger_with_audit(
target_history,
execution_history,
valuation_history,
initial_cash,
config,
)
return FactorBacktestResult(
factor_scores=factor_snapshot,
execution_prices=execution_snapshot,
valuation_prices=valuation_snapshot,
schedule=schedule,
execution_price_field=execution_field,
valuation_price_field=valuation_field,
execution=execution,
)