Files
claude_dev a2fb7b4dbb
CI/CD / test (push) Successful in 22s
CI/CD / nas-deploy (push) Successful in 12s
CI/CD / nas-verify (push) Successful in 8s
fix(portfolio): 周报双修——回测基准认 equity 列(生产四文件全被 balance-only 拒)+daily_returns 同日多行取末次(多写者防覆写) [vps]
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-25 14:49:10 +08:00

214 lines
9.3 KiB
Python

# sanguo_portfolio/weekly_report.py
"""周报三指标——流水线 spec §4.5 决议 L 六件之 3(空白新建).
TE=实际日收益 vs 回测基准日收益差 std×√252(年化);fillRate=已成交/总委托
(拒单计入分母——paper_trades 一表两态,rejected=1 即拒单);vsBacktest=周收益差
(基准=registry.backtest_run 挂的回测 run equity 序列,按日对齐).
跑在数据所在机(VPS 生产),随每日 eod 后幂等重算滚动 4 周窗;
CLI 与 APScheduler(Task 5 注册 20:45)两种触发,同一入口 run_daily_recompute.
口径注记(首年校准点): 组合实走落库(portfolio_paper.save_trade)不带拒单,
CTA PaperEngine 撮合路径才产生拒单行——分母=已落库委托行数,不补造拒单.
"""
from __future__ import annotations
import argparse
import logging
import os
import sqlite3
import statistics
from datetime import date, datetime, timedelta
from typing import Any
TRADING_DAYS = 252
ROLLING_WEEKS = 4
logger = logging.getLogger(__name__)
def _week_start(d: date) -> date:
"""该日所在 ISO 周的周一."""
return d - timedelta(days=d.weekday())
def daily_returns(series: list[tuple[str, float]]) -> dict[str, float]:
"""[(date, equity)] → {date: 当日收益率};同日多行取末次写入(多写者日末终值),
首日无前值跳过,前值 0 防除零.平盘日(收益 0.0)照记不剔除——剔除会低估 TE(M5 终审裁定).
"""
last: dict[str, float] = {}
for d, eq in series:
last[d] = eq
out: dict[str, float] = {}
prev: float | None = None
for d in sorted(last):
eq = last[d]
if prev is not None and prev != 0:
out[d] = eq / prev - 1.0
prev = eq
return out
def load_actual_daily_returns(db_path: str, account_id: int,
start: str, end: str) -> dict[str, float]:
"""paper_daily_balance.total_equity 按日 diff(Explore 实证:收盘后每账户每日一行)."""
with sqlite3.connect(db_path) as conn:
rows = conn.execute(
"SELECT date, total_equity FROM paper_daily_balance "
"WHERE account_id=? AND date>=? AND date<=? ORDER BY date",
(account_id, start, end)).fetchall()
return daily_returns([(r[0], float(r[1])) for r in rows if r[1] is not None])
def load_backtest_daily_returns(db_path: str,
task_id: str) -> dict[str, float] | None:
"""回测基准:backtest_stats 最新行 → equity_path JSON(date/balance|equity,orient=records)."""
with sqlite3.connect(db_path) as conn:
row = conn.execute(
"SELECT equity_path FROM backtest_stats WHERE task_id=? "
"ORDER BY id DESC LIMIT 1", (task_id,)).fetchone()
if row is None or not row[0] or not os.path.exists(row[0]):
return None
try:
import pandas as pd
ec = pd.read_json(row[0], orient="records")
col = "balance" if "balance" in ec.columns else (
"equity" if "equity" in ec.columns else None)
if col is None:
return None
dates = pd.to_datetime(ec["date"]).dt.strftime("%Y-%m-%d")
series = list(zip(dates.tolist(),
pd.to_numeric(ec[col], errors="coerce").tolist()))
series = [(d, float(v)) for d, v in series if v == v] # NaN 剔除
return daily_returns(series) if series else None
except Exception as exc:
logger.warning(
"load_backtest_daily_returns 解析失败 task_id=%s path=%s: %s",
task_id, row[0], exc)
return None
def compute_tracking_error(actual: dict[str, float],
baseline: dict[str, float]) -> float | None:
"""按日对齐差值样本 std×√252;<2 天不判(None)."""
common = sorted(set(actual) & set(baseline))
if len(common) < 2:
return None
diffs = [actual[d] - baseline[d] for d in common]
return statistics.stdev(diffs) * (TRADING_DAYS ** 0.5)
def compute_fill_rate(db_path: str, account_id: int,
week_start: str, week_end: str) -> float | None:
"""已成交/总委托;拒单(rejected=1)计入分母;零委托→None(不冒充 0 或 1)."""
with sqlite3.connect(db_path) as conn:
row = conn.execute(
"SELECT COUNT(*), SUM(CASE WHEN rejected=0 THEN 1 ELSE 0 END) "
"FROM paper_trades WHERE account_id=? AND bar_date>=? AND bar_date<=?",
(account_id, week_start, week_end)).fetchone()
total, filled = row[0], row[1] or 0
return filled / total if total else None
def week_return(returns: dict[str, float],
week_start: str, week_end: str) -> float | None:
"""周内复利收益;窗内无数据→None."""
rs = [v for d, v in sorted(returns.items()) if week_start <= d <= week_end]
if not rs:
return None
prod = 1.0
for r in rs:
prod *= 1.0 + r
return prod - 1.0
def build_weekly_reports(as_of: str, db_path: str, pipeline_db_path: str,
registry_path: str) -> list[dict[str, Any]]:
"""幂等重算:每策略最近 ROLLING_WEEKS 个 ISO 周各落一行(upsert).
只算 stage∈{paper,shadow,live} 且挂了 paper_account_id 的档案;
无基准(backtest_run 空/文件缺)→ te/vs 为 None,梯子页显示「无基准」.
"""
from sanguo_portfolio import pipeline_store
from sanguo_portfolio.strategy_registry import load_registry
reg = load_registry(registry_path)
today = date.fromisoformat(as_of)
rows: list[dict[str, Any]] = []
now = datetime.now().isoformat(timespec="seconds")
for name in sorted(reg["strategies"]):
entry = reg["strategies"][name]
if entry.get("stage") not in ("paper", "shadow", "live"):
continue
account_id = entry.get("paper_account_id")
if not account_id:
continue
baseline = (load_backtest_daily_returns(db_path, entry["backtest_run"])
if entry.get("backtest_run") else None)
for i in range(ROLLING_WEEKS):
ws = _week_start(today) - timedelta(weeks=i)
we = ws + timedelta(days=6)
ws_s, we_s = ws.isoformat(), we.isoformat()
end_eff = min(we, today).isoformat()
actual_w = load_actual_daily_returns(db_path, account_id, ws_s, end_eff)
if not actual_w:
continue # 该周无数据(未起跑/停牌周)不落行
fill = compute_fill_rate(db_path, account_id, ws_s, we_s)
vs = None
if baseline:
a = week_return(actual_w, ws_s, end_eff)
b = week_return(baseline, ws_s, end_eff)
vs = a - b if (a is not None and b is not None) else None
# TE=截至该周的滚动 ROLLING_WEEKS 窗(按日对齐差 std 年化)
win_start = (ws - timedelta(weeks=ROLLING_WEEKS - 1)).isoformat()
te = None
if baseline:
a_full = load_actual_daily_returns(db_path, account_id,
win_start, end_eff)
te = compute_tracking_error(a_full, baseline)
row = {"strategy_id": name, "week_start": ws_s,
"te_annual": te, "fill_rate": fill, "vs_backtest": vs,
"weeks_counted": len(actual_w), "computed_at": now}
pipeline_store.upsert_weekly(pipeline_db_path, row)
rows.append(row)
return rows
def run_daily_recompute(db_path: str, as_of: str | None = None) -> int:
"""eod 后入口(20:45 job/CLI 同源):今日 as_of 全量幂等重算.
as_of 可选注入(测试钉日期用);缺省今日——生产两调用方均不传,行为不变.
"""
from sanguo_portfolio import pipeline_store, strategy_registry
registry = os.environ.get("SANGUO_STRATEGY_REGISTRY",
os.path.join("data", "strategy_registry.yaml"))
strategy_registry.ensure_runtime_registry(registry)
as_of = as_of or date.today().isoformat()
rows = build_weekly_reports(as_of, db_path,
pipeline_store.default_pipeline_db_path(),
registry)
print(f"[weekly-report] {as_of} 重算 {len(rows)} 行(幂等 upsert)")
return len(rows)
def main(argv: list[str] | None = None) -> int:
"""CLI: python -m sanguo_portfolio.weekly_report --as-of 2026-10-09."""
ap = argparse.ArgumentParser(description="周报三指标幂等重算(决议 L 六件之 3)")
ap.add_argument("--as-of", default=None, help="缺省今日")
ap.add_argument("--db", default=None,
help="主库(SANGUO_DB_PATH 同一 db;缺省取 env 或 data/backtest_results.db)")
args = ap.parse_args(argv)
db_path = args.db or os.environ.get(
"SANGUO_DB_PATH", os.path.join("data", "backtest_results.db"))
from sanguo_portfolio import pipeline_store, strategy_registry
registry = os.environ.get("SANGUO_STRATEGY_REGISTRY",
os.path.join("data", "strategy_registry.yaml"))
strategy_registry.ensure_runtime_registry(registry)
as_of = args.as_of or date.today().isoformat()
rows = build_weekly_reports(as_of, db_path,
pipeline_store.default_pipeline_db_path(),
registry)
print(f"[weekly-report] {as_of} 落 {len(rows)} 行")
return 0
if __name__ == "__main__":
raise SystemExit(main())