214 lines
9.3 KiB
Python
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())
|