mirror of
https://ghfast.top/https://github.com/aeroxw/easy-tdx.git
synced 2026-09-12 15:44:15 +08:00
docs: add v1.13.0 portfolio management implementation plan
This commit is contained in:
@@ -0,0 +1,956 @@
|
|||||||
|
# v1.13.0 组合管理实施计划
|
||||||
|
|
||||||
|
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development or superpowers:executing-plans.
|
||||||
|
|
||||||
|
**Goal:** 实现 portfolio/ 模块(权重优化器 + 风险模型 + 再平衡引擎),使 easy-tdx 具备因子选股→组合构建→绩效评估的完整闭环。
|
||||||
|
|
||||||
|
**Architecture:** 分三层:optimizer.py(权重分配)→ risk.py(风险度量)→ rebalance.py(多期调仓回测)。每层可独立使用,RebalanceEngine 串联三者。依赖因子模块的 FactorEngine.compute_cross_section() 作为截面得分来源。
|
||||||
|
|
||||||
|
**Tech Stack:** Python 3.10+, numpy, pandas, scipy(可选,MeanVariance 降级为等权)
|
||||||
|
|
||||||
|
**Design Spec:** `docs/superpowers/specs/2026-06-12-quantitative-factor-engine-design.md` Section 4
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## File Structure
|
||||||
|
|
||||||
|
| Action | Path | Responsibility |
|
||||||
|
|--------|------|----------------|
|
||||||
|
| Create | `src/easy_tdx/portfolio/__init__.py` | 公开 API |
|
||||||
|
| Create | `src/easy_tdx/portfolio/types.py` | PortfolioState, RebalanceResult |
|
||||||
|
| Create | `src/easy_tdx/portfolio/optimizer.py` | WeightOptimizer ABC + 4 个优化器 + 注册表 |
|
||||||
|
| Create | `src/easy_tdx/portfolio/risk.py` | RiskModel (协方差/组合风险) |
|
||||||
|
| Create | `src/easy_tdx/portfolio/rebalance.py` | RebalanceEngine (多期调仓) |
|
||||||
|
| Modify | `src/easy_tdx/cli/cmd_factor.py` | 添加 portfolio 子命令 |
|
||||||
|
| Modify | `pyproject.toml` | bump → 1.13.0 |
|
||||||
|
| Create | `tests/unit/test_portfolio_optimizer.py` | 优化器测试 |
|
||||||
|
| Create | `tests/unit/test_portfolio_risk.py` | 风险模型测试 |
|
||||||
|
| Create | `tests/unit/test_portfolio_rebalance.py` | 再平衡集成测试 |
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
### Task 1: portfolio/types.py + optimizer.py
|
||||||
|
|
||||||
|
**Files:**
|
||||||
|
- Create: `src/easy_tdx/portfolio/types.py`
|
||||||
|
- Create: `src/easy_tdx/portfolio/optimizer.py`
|
||||||
|
- Create: `src/easy_tdx/portfolio/__init__.py`(初始版本)
|
||||||
|
- Test: `tests/unit/test_portfolio_optimizer.py`
|
||||||
|
|
||||||
|
#### types.py
|
||||||
|
|
||||||
|
```python
|
||||||
|
"""组合管理数据结构。"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from dataclasses import dataclass
|
||||||
|
|
||||||
|
import pandas as pd
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class PortfolioState:
|
||||||
|
"""组合状态快照。"""
|
||||||
|
|
||||||
|
date: int
|
||||||
|
weights: dict[str, float]
|
||||||
|
holdings: dict[str, float]
|
||||||
|
cash: float
|
||||||
|
total_value: float
|
||||||
|
positions_count: int
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class RebalanceResult:
|
||||||
|
"""再平衡结果。"""
|
||||||
|
|
||||||
|
rebalance_dates: list[int]
|
||||||
|
states: list[PortfolioState]
|
||||||
|
trades: pd.DataFrame
|
||||||
|
equity_curve: pd.DataFrame
|
||||||
|
performance: dict[str, float]
|
||||||
|
```
|
||||||
|
|
||||||
|
#### optimizer.py
|
||||||
|
|
||||||
|
```python
|
||||||
|
"""权重优化器。"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from abc import ABC, abstractmethod
|
||||||
|
|
||||||
|
import numpy as np
|
||||||
|
import pandas as pd
|
||||||
|
|
||||||
|
|
||||||
|
class WeightOptimizer(ABC):
|
||||||
|
"""权重优化器基类。"""
|
||||||
|
|
||||||
|
@abstractmethod
|
||||||
|
def optimize(
|
||||||
|
self,
|
||||||
|
factor_scores: pd.DataFrame, # columns: [code, score]
|
||||||
|
n_stocks: int = 50,
|
||||||
|
**kwargs: object,
|
||||||
|
) -> dict[str, float]:
|
||||||
|
"""返回 {code: weight},权重和为 1.0。"""
|
||||||
|
...
|
||||||
|
|
||||||
|
|
||||||
|
_OPTIMIZER_REGISTRY: dict[str, type[WeightOptimizer]] = {}
|
||||||
|
|
||||||
|
|
||||||
|
def register_optimizer(name: str) -> type[WeightOptimizer]:
|
||||||
|
"""注册优化器。"""
|
||||||
|
def wrapper(cls: type[WeightOptimizer]) -> type[WeightOptimizer]:
|
||||||
|
_OPTIMIZER_REGISTRY[name] = cls
|
||||||
|
return cls
|
||||||
|
return wrapper
|
||||||
|
|
||||||
|
|
||||||
|
def get_optimizer(name: str) -> WeightOptimizer:
|
||||||
|
"""按名称获取优化器实例。"""
|
||||||
|
if name not in _OPTIMIZER_REGISTRY:
|
||||||
|
raise ValueError(f"未知优化器: {name!r}。可用: {sorted(_OPTIMIZER_REGISTRY.keys())}")
|
||||||
|
return _OPTIMIZER_REGISTRY[name]()
|
||||||
|
|
||||||
|
|
||||||
|
@register_optimizer("equal")
|
||||||
|
class EqualWeightOptimizer(WeightOptimizer):
|
||||||
|
"""等权 — 取 top-N 等权分配。"""
|
||||||
|
|
||||||
|
def optimize(
|
||||||
|
self,
|
||||||
|
factor_scores: pd.DataFrame,
|
||||||
|
n_stocks: int = 50,
|
||||||
|
**kwargs: object,
|
||||||
|
) -> dict[str, float]:
|
||||||
|
top = factor_scores.nlargest(n_stocks, "score")
|
||||||
|
if len(top) == 0:
|
||||||
|
return {}
|
||||||
|
w = 1.0 / len(top)
|
||||||
|
return {row["code"]: w for _, row in top.iterrows()}
|
||||||
|
|
||||||
|
|
||||||
|
@register_optimizer("factor_weighted")
|
||||||
|
class FactorWeightedOptimizer(WeightOptimizer):
|
||||||
|
"""因子加权 — 按因子得分加权。"""
|
||||||
|
|
||||||
|
def optimize(
|
||||||
|
self,
|
||||||
|
factor_scores: pd.DataFrame,
|
||||||
|
n_stocks: int = 50,
|
||||||
|
**kwargs: object,
|
||||||
|
) -> dict[str, float]:
|
||||||
|
top = factor_scores.nlargest(n_stocks, "score")
|
||||||
|
if len(top) == 0:
|
||||||
|
return {}
|
||||||
|
scores = top["score"].to_numpy(dtype=np.float64)
|
||||||
|
# 截断极端值:top/bottom 5% 拉回到分位数
|
||||||
|
if len(scores) > 10:
|
||||||
|
q95 = np.percentile(scores, 95)
|
||||||
|
q05 = np.percentile(scores, 5)
|
||||||
|
scores = np.clip(scores, q05, q95)
|
||||||
|
# 确保正值
|
||||||
|
scores = scores - scores.min() + 1e-8
|
||||||
|
total = scores.sum()
|
||||||
|
if total == 0:
|
||||||
|
w = 1.0 / len(top)
|
||||||
|
return {row["code"]: w for _, row in top.iterrows()}
|
||||||
|
weights = scores / total
|
||||||
|
return {row["code"]: float(weights[i]) for i, (_, row) in enumerate(top.iterrows())}
|
||||||
|
|
||||||
|
|
||||||
|
@register_optimizer("risk_parity")
|
||||||
|
class RiskParityOptimizer(WeightOptimizer):
|
||||||
|
"""风险平价 — 每只股票贡献相等风险。需要 volatility 列。"""
|
||||||
|
|
||||||
|
def optimize(
|
||||||
|
self,
|
||||||
|
factor_scores: pd.DataFrame,
|
||||||
|
n_stocks: int = 50,
|
||||||
|
**kwargs: object,
|
||||||
|
) -> dict[str, float]:
|
||||||
|
top = factor_scores.nlargest(n_stocks, "score")
|
||||||
|
if len(top) == 0:
|
||||||
|
return {}
|
||||||
|
|
||||||
|
# 用 score 的倒数作为波动率代理(score 越高 = 波动率越低 = 权重越大)
|
||||||
|
# 或如果提供了 volatility 列,直接使用
|
||||||
|
if "volatility" in top.columns:
|
||||||
|
vol = top["volatility"].to_numpy(dtype=np.float64)
|
||||||
|
else:
|
||||||
|
# 用 score 均值/绝对score 作为稳定性代理
|
||||||
|
scores = top["score"].abs().to_numpy(dtype=np.float64)
|
||||||
|
vol = 1.0 / (scores + 1e-8)
|
||||||
|
|
||||||
|
vol = np.maximum(vol, 1e-8)
|
||||||
|
inv_vol = 1.0 / vol
|
||||||
|
total = inv_vol.sum()
|
||||||
|
weights = inv_vol / total
|
||||||
|
|
||||||
|
return {row["code"]: float(weights[i]) for i, (_, row) in enumerate(top.iterrows())}
|
||||||
|
|
||||||
|
|
||||||
|
@register_optimizer("mean_variance")
|
||||||
|
class MeanVarianceOptimizer(WeightOptimizer):
|
||||||
|
"""均值方差优化 — Markowitz 模型(可选 scipy)。"""
|
||||||
|
|
||||||
|
def optimize(
|
||||||
|
self,
|
||||||
|
factor_scores: pd.DataFrame,
|
||||||
|
n_stocks: int = 50,
|
||||||
|
**kwargs: object,
|
||||||
|
) -> dict[str, float]:
|
||||||
|
# 尝试使用 scipy,降级到等权
|
||||||
|
try:
|
||||||
|
from scipy.optimize import minimize
|
||||||
|
return self._optimize_with_scipy(factor_scores, n_stocks, minimize)
|
||||||
|
except ImportError:
|
||||||
|
# scipy 不可用,降级为等权
|
||||||
|
fallback = EqualWeightOptimizer()
|
||||||
|
return fallback.optimize(factor_scores, n_stocks)
|
||||||
|
|
||||||
|
def _optimize_with_scipy(
|
||||||
|
self,
|
||||||
|
factor_scores: pd.DataFrame,
|
||||||
|
n_stocks: int,
|
||||||
|
minimize: object,
|
||||||
|
) -> dict[str, float]:
|
||||||
|
from scipy.optimize import minimize as _minimize # type: ignore[import]
|
||||||
|
|
||||||
|
top = factor_scores.nlargest(n_stocks, "score")
|
||||||
|
if len(top) == 0:
|
||||||
|
return {}
|
||||||
|
|
||||||
|
n = len(top)
|
||||||
|
scores = top["score"].to_numpy(dtype=np.float64)
|
||||||
|
|
||||||
|
# 用 score 构造简化协方差矩阵(对角矩阵,方差 = 1/score²)
|
||||||
|
variances = 1.0 / (np.abs(scores) + 1e-8) ** 2
|
||||||
|
cov = np.diag(variances)
|
||||||
|
|
||||||
|
# 目标:最小化 w' Σ w,约束 sum(w)=1, 0<=w<=0.1
|
||||||
|
def objective(w: np.ndarray) -> float:
|
||||||
|
return float(w @ cov @ w)
|
||||||
|
|
||||||
|
constraints = {"type": "eq", "fun": lambda w: float(np.sum(w) - 1.0)}
|
||||||
|
bounds = [(0.0, 0.1)] * n
|
||||||
|
x0 = np.ones(n) / n
|
||||||
|
|
||||||
|
result = _minimize(objective, x0, method="SLSQP", bounds=bounds, constraints=constraints)
|
||||||
|
|
||||||
|
if result.success:
|
||||||
|
weights = result.x
|
||||||
|
else:
|
||||||
|
weights = np.ones(n) / n
|
||||||
|
|
||||||
|
return {row["code"]: float(weights[i]) for i, (_, row) in enumerate(top.iterrows())}
|
||||||
|
```
|
||||||
|
|
||||||
|
#### 测试
|
||||||
|
|
||||||
|
```python
|
||||||
|
# tests/unit/test_portfolio_optimizer.py
|
||||||
|
"""Test portfolio optimizers."""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import numpy as np
|
||||||
|
import pandas as pd
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from easy_tdx.portfolio.optimizer import (
|
||||||
|
EqualWeightOptimizer,
|
||||||
|
FactorWeightedOptimizer,
|
||||||
|
RiskParityOptimizer,
|
||||||
|
get_optimizer,
|
||||||
|
register_optimizer,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _make_scores(n: int = 20, seed: int = 42) -> pd.DataFrame:
|
||||||
|
rng = np.random.default_rng(seed)
|
||||||
|
return pd.DataFrame({
|
||||||
|
"code": [f"{i:06d}" for i in range(n)],
|
||||||
|
"score": rng.normal(0.02, 0.05, n),
|
||||||
|
})
|
||||||
|
|
||||||
|
|
||||||
|
class TestEqualWeight:
|
||||||
|
def test_weights_sum_to_one(self):
|
||||||
|
opt = EqualWeightOptimizer()
|
||||||
|
w = opt.optimize(_make_scores(), n_stocks=10)
|
||||||
|
assert abs(sum(w.values()) - 1.0) < 1e-8
|
||||||
|
|
||||||
|
def test_n_stocks_selected(self):
|
||||||
|
opt = EqualWeightOptimizer()
|
||||||
|
w = opt.optimize(_make_scores(), n_stocks=5)
|
||||||
|
assert len(w) == 5
|
||||||
|
|
||||||
|
def test_all_equal(self):
|
||||||
|
opt = EqualWeightOptimizer()
|
||||||
|
w = opt.optimize(_make_scores(), n_stocks=10)
|
||||||
|
values = list(w.values())
|
||||||
|
assert all(abs(v - values[0]) < 1e-8 for v in values)
|
||||||
|
|
||||||
|
def test_empty_input(self):
|
||||||
|
opt = EqualWeightOptimizer()
|
||||||
|
w = opt.optimize(pd.DataFrame(columns=["code", "score"]), n_stocks=5)
|
||||||
|
assert len(w) == 0
|
||||||
|
|
||||||
|
|
||||||
|
class TestFactorWeighted:
|
||||||
|
def test_weights_sum_to_one(self):
|
||||||
|
opt = FactorWeightedOptimizer()
|
||||||
|
w = opt.optimize(_make_scores(), n_stocks=10)
|
||||||
|
assert abs(sum(w.values()) - 1.0) < 1e-6
|
||||||
|
|
||||||
|
def test_higher_score_higher_weight(self):
|
||||||
|
scores = pd.DataFrame({"code": ["A", "B", "C"], "score": [3.0, 2.0, 1.0]})
|
||||||
|
opt = FactorWeightedOptimizer()
|
||||||
|
w = opt.optimize(scores, n_stocks=3)
|
||||||
|
assert w["A"] > w["C"]
|
||||||
|
|
||||||
|
|
||||||
|
class TestRiskParity:
|
||||||
|
def test_weights_sum_to_one(self):
|
||||||
|
opt = RiskParityOptimizer()
|
||||||
|
w = opt.optimize(_make_scores(), n_stocks=10)
|
||||||
|
assert abs(sum(w.values()) - 1.0) < 1e-6
|
||||||
|
|
||||||
|
def test_with_volatility_column(self):
|
||||||
|
scores = pd.DataFrame({
|
||||||
|
"code": ["A", "B", "C"],
|
||||||
|
"score": [1.0, 1.0, 1.0],
|
||||||
|
"volatility": [0.1, 0.2, 0.4],
|
||||||
|
})
|
||||||
|
opt = RiskParityOptimizer()
|
||||||
|
w = opt.optimize(scores, n_stocks=3)
|
||||||
|
# 低波动率 → 更高权重
|
||||||
|
assert w["A"] > w["C"]
|
||||||
|
|
||||||
|
|
||||||
|
class TestRegistry:
|
||||||
|
def test_get_optimizer(self):
|
||||||
|
opt = get_optimizer("equal")
|
||||||
|
assert isinstance(opt, EqualWeightOptimizer)
|
||||||
|
|
||||||
|
def test_unknown_raises(self):
|
||||||
|
with pytest.raises(ValueError, match="未知优化器"):
|
||||||
|
get_optimizer("nonexistent")
|
||||||
|
```
|
||||||
|
|
||||||
|
#### __init__.py(初始)
|
||||||
|
|
||||||
|
```python
|
||||||
|
"""组合管理模块。"""
|
||||||
|
```
|
||||||
|
|
||||||
|
提交: `feat(portfolio): add optimizer module with 4 strategies (equal/factor_weighted/risk_parity/mean_variance)`
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
### Task 2: portfolio/risk.py
|
||||||
|
|
||||||
|
```python
|
||||||
|
# src/easy_tdx/portfolio/risk.py
|
||||||
|
"""简化风险模型。"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import numpy as np
|
||||||
|
import pandas as pd
|
||||||
|
|
||||||
|
|
||||||
|
class RiskModel:
|
||||||
|
"""简化风险模型 — A 股够用。"""
|
||||||
|
|
||||||
|
def estimate_covariance(
|
||||||
|
self,
|
||||||
|
returns: pd.DataFrame,
|
||||||
|
method: str = "shrinkage",
|
||||||
|
shrinkage_intensity: float = 0.5,
|
||||||
|
window: int = 60,
|
||||||
|
) -> pd.DataFrame:
|
||||||
|
"""协方差矩阵估计。
|
||||||
|
|
||||||
|
Args:
|
||||||
|
returns: columns=code, index=date 的收益率 DataFrame。
|
||||||
|
method: "shrinkage" | "sample"
|
||||||
|
"""
|
||||||
|
if len(returns) < 2:
|
||||||
|
codes = returns.columns.tolist() if len(returns.columns) > 0 else []
|
||||||
|
return pd.DataFrame(np.eye(len(codes)), index=codes, columns=codes)
|
||||||
|
|
||||||
|
# 只用最近 window 行
|
||||||
|
if len(returns) > window:
|
||||||
|
returns = returns.iloc[-window:]
|
||||||
|
|
||||||
|
sample_cov = returns.cov()
|
||||||
|
|
||||||
|
if method == "shrinkage":
|
||||||
|
# Ledoit-Wolf 简化版:收缩到对角矩阵
|
||||||
|
target = pd.DataFrame(
|
||||||
|
np.diag(np.diag(sample_cov.to_numpy())),
|
||||||
|
index=sample_cov.index,
|
||||||
|
columns=sample_cov.columns,
|
||||||
|
)
|
||||||
|
shrunk = (1 - shrinkage_intensity) * sample_cov + shrinkage_intensity * target
|
||||||
|
return shrunk
|
||||||
|
|
||||||
|
return sample_cov
|
||||||
|
|
||||||
|
def portfolio_risk(
|
||||||
|
self,
|
||||||
|
weights: dict[str, float],
|
||||||
|
cov_matrix: pd.DataFrame,
|
||||||
|
) -> dict[str, float]:
|
||||||
|
"""组合风险指标。
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
total_volatility, max_risk_contribution, n_positions
|
||||||
|
"""
|
||||||
|
codes = [c for c in weights if c in cov_matrix.columns]
|
||||||
|
if not codes:
|
||||||
|
return {"total_volatility": 0.0, "max_risk_contribution": 0.0, "n_positions": 0}
|
||||||
|
|
||||||
|
w = np.array([weights[c] for c in codes])
|
||||||
|
cov_sub = cov_matrix.loc[codes, codes].to_numpy()
|
||||||
|
|
||||||
|
# 总波动率(年化)
|
||||||
|
var = float(w @ cov_sub @ w)
|
||||||
|
total_vol = np.sqrt(max(0, var)) * np.sqrt(252)
|
||||||
|
|
||||||
|
# 边际风险贡献
|
||||||
|
marginal = cov_sub @ w
|
||||||
|
risk_contrib = np.abs(w * marginal)
|
||||||
|
total_rc = risk_contrib.sum()
|
||||||
|
max_rc = float(risk_contrib.max() / total_rc) if total_rc > 0 else 0.0
|
||||||
|
|
||||||
|
return {
|
||||||
|
"total_volatility": total_vol,
|
||||||
|
"max_risk_contribution": max_rc,
|
||||||
|
"n_positions": len(codes),
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
测试:
|
||||||
|
|
||||||
|
```python
|
||||||
|
# tests/unit/test_portfolio_risk.py
|
||||||
|
"""Test RiskModel."""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import numpy as np
|
||||||
|
import pandas as pd
|
||||||
|
|
||||||
|
from easy_tdx.portfolio.risk import RiskModel
|
||||||
|
|
||||||
|
|
||||||
|
def _make_returns(n_dates: int = 100, n_stocks: int = 5, seed: int = 42) -> pd.DataFrame:
|
||||||
|
rng = np.random.default_rng(seed)
|
||||||
|
codes = [f"{i:06d}" for i in range(n_stocks)]
|
||||||
|
data = rng.normal(0.001, 0.02, (n_dates, n_stocks))
|
||||||
|
return pd.DataFrame(data, columns=codes)
|
||||||
|
|
||||||
|
|
||||||
|
class TestCovarianceEstimation:
|
||||||
|
def test_shape(self):
|
||||||
|
rm = RiskModel()
|
||||||
|
ret = _make_returns()
|
||||||
|
cov = rm.estimate_covariance(ret)
|
||||||
|
assert cov.shape == (5, 5)
|
||||||
|
|
||||||
|
def test_symmetric(self):
|
||||||
|
rm = RiskModel()
|
||||||
|
ret = _make_returns()
|
||||||
|
cov = rm.estimate_covariance(ret)
|
||||||
|
assert np.allclose(cov.to_numpy(), cov.to_numpy().T)
|
||||||
|
|
||||||
|
def test_shrinkage_vs_sample(self):
|
||||||
|
rm = RiskModel()
|
||||||
|
ret = _make_returns()
|
||||||
|
shrunk = rm.estimate_covariance(ret, method="shrinkage")
|
||||||
|
sample = rm.estimate_covariance(ret, method="sample")
|
||||||
|
# 收缩矩阵的非对角元素应更小
|
||||||
|
off_diag_shrunk = shrunk.values[~np.eye(len(shrunk), dtype=bool)]
|
||||||
|
off_diag_sample = sample.values[~np.eye(len(sample), dtype=bool)]
|
||||||
|
assert np.abs(off_diag_shrunk).mean() <= np.abs(off_diag_sample).mean()
|
||||||
|
|
||||||
|
|
||||||
|
class TestPortfolioRisk:
|
||||||
|
def test_total_volatility(self):
|
||||||
|
rm = RiskModel()
|
||||||
|
ret = _make_returns()
|
||||||
|
cov = rm.estimate_covariance(ret)
|
||||||
|
weights = {"000000": 0.5, "000001": 0.5}
|
||||||
|
risk = rm.portfolio_risk(weights, cov)
|
||||||
|
assert risk["total_volatility"] > 0
|
||||||
|
assert risk["n_positions"] == 2
|
||||||
|
|
||||||
|
def test_empty_weights(self):
|
||||||
|
rm = RiskModel()
|
||||||
|
cov = pd.DataFrame()
|
||||||
|
risk = rm.portfolio_risk({}, cov)
|
||||||
|
assert risk["total_volatility"] == 0.0
|
||||||
|
```
|
||||||
|
|
||||||
|
提交: `feat(portfolio): add RiskModel with shrinkage covariance and portfolio risk decomposition`
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
### Task 3: portfolio/rebalance.py — 再平衡引擎
|
||||||
|
|
||||||
|
```python
|
||||||
|
# src/easy_tdx/portfolio/rebalance.py
|
||||||
|
"""多期调仓回测引擎。"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import numpy as np
|
||||||
|
import pandas as pd
|
||||||
|
|
||||||
|
from easy_tdx.factor.engine import FactorEngine
|
||||||
|
from easy_tdx.portfolio.optimizer import WeightOptimizer
|
||||||
|
from easy_tdx.portfolio.types import PortfolioState, RebalanceResult
|
||||||
|
|
||||||
|
|
||||||
|
class RebalanceEngine:
|
||||||
|
"""多期调仓回测引擎。
|
||||||
|
|
||||||
|
管道: FactorEngine → 截面得分 → Optimizer → 目标权重 → 交易成本 → 持仓跟踪
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
optimizer: WeightOptimizer,
|
||||||
|
factor_name: str = "momentum_20d",
|
||||||
|
n_stocks: int = 50,
|
||||||
|
rebalance_freq: str = "M",
|
||||||
|
commission: float = 0.0003,
|
||||||
|
slippage: float = 0.001,
|
||||||
|
cash: float = 1_000_000,
|
||||||
|
) -> None:
|
||||||
|
self._optimizer = optimizer
|
||||||
|
self._factor_name = factor_name
|
||||||
|
self._n_stocks = n_stocks
|
||||||
|
self._rebalance_freq = rebalance_freq
|
||||||
|
self._commission = commission
|
||||||
|
self._slippage = slippage
|
||||||
|
self._cash = cash
|
||||||
|
|
||||||
|
def _get_rebalance_dates(self, dates: pd.DatetimeIndex) -> list[int]:
|
||||||
|
"""根据频率确定调仓日期。"""
|
||||||
|
freq_map = {"W": "W-MON", "M": "ME", "Q": "QE"}
|
||||||
|
freq = freq_map.get(self._rebalance_freq, "ME")
|
||||||
|
|
||||||
|
series = pd.Series(dates, index=dates)
|
||||||
|
grouped = series.groupby(series.dt.to_period(freq))
|
||||||
|
rebalance_dates = [group.iloc[-1] for _, group in grouped if len(group) > 0]
|
||||||
|
return [int(d.strftime("%Y%m%d")) for d in rebalance_dates]
|
||||||
|
|
||||||
|
def run(
|
||||||
|
self,
|
||||||
|
data: dict[str, pd.DataFrame],
|
||||||
|
start_date: int | None = None,
|
||||||
|
end_date: int | None = None,
|
||||||
|
) -> RebalanceResult:
|
||||||
|
"""执行多期回测。"""
|
||||||
|
if not data:
|
||||||
|
return self._empty_result()
|
||||||
|
|
||||||
|
factor_engine = FactorEngine()
|
||||||
|
|
||||||
|
# 收集所有日期
|
||||||
|
all_dates: list[pd.Timestamp] = []
|
||||||
|
for df in data.values():
|
||||||
|
if "datetime" in df.columns:
|
||||||
|
all_dates.extend(df["datetime"].tolist())
|
||||||
|
if not all_dates:
|
||||||
|
return self._empty_result()
|
||||||
|
|
||||||
|
all_dates = sorted(set(all_dates))
|
||||||
|
date_ints = [int(d.strftime("%Y%m%d")) for d in all_dates]
|
||||||
|
|
||||||
|
# 日期范围过滤
|
||||||
|
if start_date:
|
||||||
|
all_dates = [d for d, di in zip(all_dates, date_ints) if di >= start_date]
|
||||||
|
if end_date:
|
||||||
|
all_dates = [d for d in all_dates if int(d.strftime("%Y%m%d")) <= end_date]
|
||||||
|
|
||||||
|
if not all_dates:
|
||||||
|
return self._empty_result()
|
||||||
|
|
||||||
|
# 确定调仓日期
|
||||||
|
dt_index = pd.DatetimeIndex(all_dates)
|
||||||
|
rebalance_dates = self._get_rebalance_dates(dt_index)
|
||||||
|
rebalance_set = set(rebalance_dates)
|
||||||
|
|
||||||
|
# 模拟
|
||||||
|
cash = self._cash
|
||||||
|
holdings: dict[str, float] = {} # code -> shares
|
||||||
|
states: list[PortfolioState] = []
|
||||||
|
trades_list: list[dict[str, object]] = []
|
||||||
|
equity_records: list[dict[str, object]] = []
|
||||||
|
|
||||||
|
for dt in all_dates:
|
||||||
|
date_int = int(dt.strftime("%Y%m%d"))
|
||||||
|
is_rebalance = date_int in rebalance_set
|
||||||
|
|
||||||
|
# 获取当天各股票价格
|
||||||
|
prices: dict[str, float] = {}
|
||||||
|
for code, df in data.items():
|
||||||
|
if "datetime" in df.columns:
|
||||||
|
row = df[df["datetime"] == dt]
|
||||||
|
if not row.empty:
|
||||||
|
prices[code] = float(row["close"].iloc[0])
|
||||||
|
|
||||||
|
# 计算持仓市值
|
||||||
|
position_value = sum(holdings.get(c, 0) * prices.get(c, 0) for c in holdings)
|
||||||
|
total_value = cash + position_value
|
||||||
|
|
||||||
|
if is_rebalance and total_value > 0:
|
||||||
|
# 计算截面因子得分
|
||||||
|
scores_df = factor_engine.compute_cross_section(
|
||||||
|
data, [self._factor_name], date=date_int
|
||||||
|
)
|
||||||
|
|
||||||
|
if not scores_df.empty and scores_df[self._factor_name].notna().any():
|
||||||
|
scores_df = scores_df.rename(columns={self._factor_name: "score"})
|
||||||
|
scores_df = scores_df[["code", "score"]].dropna(subset=["score"])
|
||||||
|
|
||||||
|
# 优化权重
|
||||||
|
target_weights = self._optimizer.optimize(scores_df, n_stocks=self._n_stocks)
|
||||||
|
else:
|
||||||
|
target_weights = {}
|
||||||
|
|
||||||
|
# 执行调仓
|
||||||
|
if target_weights:
|
||||||
|
trades_list, cash, holdings = self._rebalance(
|
||||||
|
target_weights, prices, total_value, date_int, trades_list, cash, holdings
|
||||||
|
)
|
||||||
|
|
||||||
|
# 记录状态
|
||||||
|
position_value = sum(holdings.get(c, 0) * prices.get(c, 0) for c in holdings)
|
||||||
|
total_value = cash + position_value
|
||||||
|
weights = {}
|
||||||
|
if total_value > 0:
|
||||||
|
for c in holdings:
|
||||||
|
if holdings[c] > 0 and c in prices:
|
||||||
|
weights[c] = holdings[c] * prices[c] / total_value
|
||||||
|
|
||||||
|
states.append(PortfolioState(
|
||||||
|
date=date_int,
|
||||||
|
weights=weights,
|
||||||
|
holdings=dict(holdings),
|
||||||
|
cash=cash,
|
||||||
|
total_value=total_value,
|
||||||
|
positions_count=len([s for s in holdings.values() if s > 0]),
|
||||||
|
))
|
||||||
|
|
||||||
|
equity_records.append({
|
||||||
|
"datetime": date_int,
|
||||||
|
"total": total_value,
|
||||||
|
"cash": cash,
|
||||||
|
"position_value": position_value,
|
||||||
|
})
|
||||||
|
|
||||||
|
# 构建结果
|
||||||
|
equity_curve = pd.DataFrame(equity_records)
|
||||||
|
trades_df = pd.DataFrame(trades_list) if trades_list else pd.DataFrame(
|
||||||
|
columns=["datetime", "direction", "code", "shares", "price", "cost"]
|
||||||
|
)
|
||||||
|
|
||||||
|
performance = self._compute_performance(equity_curve)
|
||||||
|
|
||||||
|
return RebalanceResult(
|
||||||
|
rebalance_dates=rebalance_dates,
|
||||||
|
states=states,
|
||||||
|
trades=trades_df,
|
||||||
|
equity_curve=equity_curve,
|
||||||
|
performance=performance,
|
||||||
|
)
|
||||||
|
|
||||||
|
def _rebalance(
|
||||||
|
self,
|
||||||
|
target_weights: dict[str, float],
|
||||||
|
prices: dict[str, float],
|
||||||
|
total_value: float,
|
||||||
|
date_int: int,
|
||||||
|
trades_list: list[dict[str, object]],
|
||||||
|
cash: float,
|
||||||
|
holdings: dict[str, float],
|
||||||
|
) -> tuple[list[dict[str, object]], float, dict[str, float]]:
|
||||||
|
"""执行一次调仓。"""
|
||||||
|
# 先卖出不在目标组合中的持仓
|
||||||
|
for code in list(holdings.keys()):
|
||||||
|
if code not in target_weights and holdings[code] > 0:
|
||||||
|
price = prices.get(code, 0)
|
||||||
|
if price > 0:
|
||||||
|
sell_value = holdings[code] * price
|
||||||
|
cost = sell_value * (self._commission + self._slippage)
|
||||||
|
cash += sell_value - cost
|
||||||
|
trades_list.append({
|
||||||
|
"datetime": date_int,
|
||||||
|
"direction": "SELL",
|
||||||
|
"code": code,
|
||||||
|
"shares": holdings[code],
|
||||||
|
"price": price,
|
||||||
|
"cost": cost,
|
||||||
|
})
|
||||||
|
del holdings[code]
|
||||||
|
|
||||||
|
# 调整目标持仓
|
||||||
|
new_holdings: dict[str, float] = {}
|
||||||
|
total_cost = 0.0
|
||||||
|
for code, weight in target_weights.items():
|
||||||
|
price = prices.get(code, 0)
|
||||||
|
if price <= 0:
|
||||||
|
continue
|
||||||
|
target_value = total_value * weight
|
||||||
|
shares = int(target_value / price / 100) * 100 # 100 股整手
|
||||||
|
if shares > 0:
|
||||||
|
new_holdings[code] = shares
|
||||||
|
trade_value = shares * price
|
||||||
|
cost = trade_value * (self._commission + self._slippage)
|
||||||
|
total_cost += trade_value + cost
|
||||||
|
trades_list.append({
|
||||||
|
"datetime": date_int,
|
||||||
|
"direction": "BUY",
|
||||||
|
"code": code,
|
||||||
|
"shares": shares,
|
||||||
|
"price": price,
|
||||||
|
"cost": cost,
|
||||||
|
})
|
||||||
|
|
||||||
|
cash = total_value - sum(
|
||||||
|
new_holdings.get(c, 0) * prices.get(c, 0) for c in new_holdings
|
||||||
|
)
|
||||||
|
holdings.clear()
|
||||||
|
holdings.update(new_holdings)
|
||||||
|
|
||||||
|
return trades_list, cash, holdings
|
||||||
|
|
||||||
|
def _compute_performance(self, equity_curve: pd.DataFrame) -> dict[str, float]:
|
||||||
|
"""计算绩效指标。"""
|
||||||
|
if len(equity_curve) < 2:
|
||||||
|
return {"total_return": 0.0, "annual_return": 0.0, "max_drawdown": 0.0, "sharpe": 0.0}
|
||||||
|
|
||||||
|
total = equity_curve["total"].to_numpy()
|
||||||
|
total_return = (total[-1] / total[0]) - 1
|
||||||
|
|
||||||
|
# 年化收益
|
||||||
|
n_days = len(total)
|
||||||
|
annual_return = (1 + total_return) ** (252 / max(n_days, 1)) - 1
|
||||||
|
|
||||||
|
# 最大回撤
|
||||||
|
peak = np.maximum.accumulate(total)
|
||||||
|
drawdown = (total - peak) / peak
|
||||||
|
max_drawdown = float(np.min(drawdown))
|
||||||
|
|
||||||
|
# 夏普比率
|
||||||
|
daily_ret = np.diff(total) / total[:-1]
|
||||||
|
daily_ret = daily_ret[~np.isnan(daily_ret)]
|
||||||
|
if len(daily_ret) > 1 and np.std(daily_ret) > 0:
|
||||||
|
sharpe = float(np.mean(daily_ret) / np.std(daily_ret) * np.sqrt(252))
|
||||||
|
else:
|
||||||
|
sharpe = 0.0
|
||||||
|
|
||||||
|
return {
|
||||||
|
"total_return": total_return,
|
||||||
|
"annual_return": annual_return,
|
||||||
|
"max_drawdown": max_drawdown,
|
||||||
|
"sharpe": sharpe,
|
||||||
|
"total_trades": len(equity_curve),
|
||||||
|
}
|
||||||
|
|
||||||
|
def _empty_result(self) -> RebalanceResult:
|
||||||
|
return RebalanceResult(
|
||||||
|
rebalance_dates=[],
|
||||||
|
states=[],
|
||||||
|
trades=pd.DataFrame(columns=["datetime", "direction", "code", "shares", "price", "cost"]),
|
||||||
|
equity_curve=pd.DataFrame(columns=["datetime", "total", "cash", "position_value"]),
|
||||||
|
performance={"total_return": 0.0, "annual_return": 0.0, "max_drawdown": 0.0, "sharpe": 0.0},
|
||||||
|
)
|
||||||
|
```
|
||||||
|
|
||||||
|
测试:
|
||||||
|
|
||||||
|
```python
|
||||||
|
# tests/unit/test_portfolio_rebalance.py
|
||||||
|
"""Test RebalanceEngine."""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import numpy as np
|
||||||
|
import pandas as pd
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from easy_tdx.portfolio.optimizer import EqualWeightOptimizer, FactorWeightedOptimizer
|
||||||
|
from easy_tdx.portfolio.rebalance import RebalanceEngine
|
||||||
|
|
||||||
|
|
||||||
|
def _make_market(n_stocks: int = 10, n_days: int = 120, seed: int = 42) -> dict[str, pd.DataFrame]:
|
||||||
|
"""模拟一个微型市场。"""
|
||||||
|
rng = np.random.default_rng(seed)
|
||||||
|
data = {}
|
||||||
|
for i in range(n_stocks):
|
||||||
|
close = 10.0 + np.cumsum(rng.normal(0.01, 0.5, n_days))
|
||||||
|
close = np.maximum(close, 1.0)
|
||||||
|
data[f"{i:06d}"] = pd.DataFrame({
|
||||||
|
"datetime": pd.date_range("2024-01-01", periods=n_days, freq="D"),
|
||||||
|
"open": close,
|
||||||
|
"high": close + 0.3,
|
||||||
|
"low": close - 0.3,
|
||||||
|
"close": close,
|
||||||
|
"vol": rng.integers(1e5, 1e7, n_days).astype(float),
|
||||||
|
"amount": close * 1e6,
|
||||||
|
})
|
||||||
|
return data
|
||||||
|
|
||||||
|
|
||||||
|
class TestRebalanceEngine:
|
||||||
|
def test_basic_run(self):
|
||||||
|
engine = RebalanceEngine(
|
||||||
|
optimizer=EqualWeightOptimizer(),
|
||||||
|
factor_name="momentum_20d",
|
||||||
|
n_stocks=5,
|
||||||
|
rebalance_freq="M",
|
||||||
|
cash=1_000_000,
|
||||||
|
)
|
||||||
|
data = _make_market()
|
||||||
|
result = engine.run(data, start_date=20240101, end_date=20240430)
|
||||||
|
|
||||||
|
assert len(result.states) > 0
|
||||||
|
assert len(result.rebalance_dates) > 0
|
||||||
|
assert len(result.equity_curve) > 0
|
||||||
|
assert "total_return" in result.performance
|
||||||
|
|
||||||
|
def test_with_factor_weighted(self):
|
||||||
|
engine = RebalanceEngine(
|
||||||
|
optimizer=FactorWeightedOptimizer(),
|
||||||
|
factor_name="momentum_20d",
|
||||||
|
n_stocks=5,
|
||||||
|
rebalance_freq="M",
|
||||||
|
)
|
||||||
|
data = _make_market()
|
||||||
|
result = engine.run(data, start_date=20240101, end_date=20240430)
|
||||||
|
assert len(result.states) > 0
|
||||||
|
|
||||||
|
def test_empty_data(self):
|
||||||
|
engine = RebalanceEngine(optimizer=EqualWeightOptimizer())
|
||||||
|
result = engine.run({})
|
||||||
|
assert result.performance["total_return"] == 0.0
|
||||||
|
|
||||||
|
def test_equity_curve_monotonic_dates(self):
|
||||||
|
engine = RebalanceEngine(optimizer=EqualWeightOptimizer(), rebalance_freq="M")
|
||||||
|
data = _make_market()
|
||||||
|
result = engine.run(data, start_date=20240101, end_date=20240430)
|
||||||
|
dates = result.equity_curve["datetime"].tolist()
|
||||||
|
assert dates == sorted(dates)
|
||||||
|
|
||||||
|
def test_trades_recorded(self):
|
||||||
|
engine = RebalanceEngine(
|
||||||
|
optimizer=EqualWeightOptimizer(),
|
||||||
|
n_stocks=3,
|
||||||
|
rebalance_freq="M",
|
||||||
|
)
|
||||||
|
data = _make_market()
|
||||||
|
result = engine.run(data, start_date=20240101, end_date=20240430)
|
||||||
|
assert len(result.trades) > 0
|
||||||
|
assert "BUY" in result.trades["direction"].values
|
||||||
|
```
|
||||||
|
|
||||||
|
提交: `feat(portfolio): add RebalanceEngine with multi-period portfolio backtesting`
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
### Task 4: portfolio/__init__.py + CLI + 版本号
|
||||||
|
|
||||||
|
#### __init__.py
|
||||||
|
|
||||||
|
```python
|
||||||
|
"""组合管理模块。"""
|
||||||
|
from easy_tdx.portfolio.optimizer import (
|
||||||
|
EqualWeightOptimizer,
|
||||||
|
FactorWeightedOptimizer,
|
||||||
|
MeanVarianceOptimizer,
|
||||||
|
RiskParityOptimizer,
|
||||||
|
WeightOptimizer,
|
||||||
|
get_optimizer,
|
||||||
|
)
|
||||||
|
from easy_tdx.portfolio.rebalance import RebalanceEngine
|
||||||
|
from easy_tdx.portfolio.risk import RiskModel
|
||||||
|
from easy_tdx.portfolio.types import PortfolioState, RebalanceResult
|
||||||
|
|
||||||
|
__all__ = [
|
||||||
|
"WeightOptimizer",
|
||||||
|
"EqualWeightOptimizer",
|
||||||
|
"FactorWeightedOptimizer",
|
||||||
|
"RiskParityOptimizer",
|
||||||
|
"MeanVarianceOptimizer",
|
||||||
|
"get_optimizer",
|
||||||
|
"RiskModel",
|
||||||
|
"RebalanceEngine",
|
||||||
|
"PortfolioState",
|
||||||
|
"RebalanceResult",
|
||||||
|
]
|
||||||
|
```
|
||||||
|
|
||||||
|
#### CLI: 在 cli/__init__.py 注册 portfolio 命令
|
||||||
|
|
||||||
|
先检查是否已有 portfolio 命令(backtest 模块有一个 portfolio 子命令)。如果有,跳过 CLI 修改,只更新 portfolio/__init__.py。
|
||||||
|
|
||||||
|
如果没有冲突,创建 src/easy_tdx/cli/cmd_portfolio.py:
|
||||||
|
|
||||||
|
```python
|
||||||
|
"""组合管理 CLI 命令。"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import click
|
||||||
|
|
||||||
|
|
||||||
|
@click.group("portfolio")
|
||||||
|
def portfolio() -> None:
|
||||||
|
"""组合管理工具。"""
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
@portfolio.command("backtest")
|
||||||
|
@click.argument("factor_name")
|
||||||
|
@click.option("--n-stocks", default=50, type=int, help="持仓数量")
|
||||||
|
@click.option("--rebalance-freq", default="M", help="调仓频率: W/M/Q")
|
||||||
|
@click.option("--optimizer", "opt_name", default="equal", help="优化器: equal/factor_weighted/risk_parity/mean_variance")
|
||||||
|
@click.option("--cash", default=1000000.0, type=float, help="初始资金")
|
||||||
|
def portfolio_backtest(factor_name: str, n_stocks: int, rebalance_freq: str, opt_name: str, cash: float) -> None:
|
||||||
|
"""运行组合回测。
|
||||||
|
|
||||||
|
示例:
|
||||||
|
|
||||||
|
easy-tdx portfolio backtest momentum_20d
|
||||||
|
|
||||||
|
easy-tdx portfolio backtest rsi_14 --n-stocks 10 --optimizer factor_weighted
|
||||||
|
"""
|
||||||
|
import json
|
||||||
|
click.echo(json.dumps({
|
||||||
|
"message": "portfolio backtest 需要行情数据,请使用 Python API",
|
||||||
|
"example": (
|
||||||
|
f"from easy_tdx.portfolio import RebalanceEngine, EqualWeightOptimizer\n"
|
||||||
|
f"from easy_tdx.factor import FactorEngine\n"
|
||||||
|
f"engine = RebalanceEngine(\n"
|
||||||
|
f" optimizer=EqualWeightOptimizer(),\n"
|
||||||
|
f" factor_name='{factor_name}',\n"
|
||||||
|
f" n_stocks={n_stocks},\n"
|
||||||
|
f" rebalance_freq='{rebalance_freq}',\n"
|
||||||
|
f" cash={cash},\n"
|
||||||
|
f")\n"
|
||||||
|
f"result = engine.run(data, start_date=20230101, end_date=20240101)\n"
|
||||||
|
f"print(f'年化收益={{result.performance[\"annual_return\"]:.2%}}')"
|
||||||
|
),
|
||||||
|
}, ensure_ascii=False, indent=2))
|
||||||
|
```
|
||||||
|
|
||||||
|
在 cli/__init__.py 中导入并注册。
|
||||||
|
|
||||||
|
#### pyproject.toml: version = "1.13.0"
|
||||||
|
|
||||||
|
#### 运行全量测试 + ruff + 提交
|
||||||
|
|
||||||
|
提交: `feat(portfolio): add portfolio module exports, CLI command, bump v1.13.0`
|
||||||
Reference in New Issue
Block a user