feat: 策略数据每日自动更新——period_rollup滑窗派生1y/2y/5y行(锚定数据末端,不覆盖真实记录)+health_monitor数据驱动跟随路由激活集合+regime_perf排序改最长周期优先(防派生行污染温区行)
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
#!/usr/bin/env python3
|
||||
"""health_monitor_daily.py — 策略健康度每日监控(v_weak + v_oversold)
|
||||
"""health_monitor_daily.py — 策略健康度每日监控(数据驱动:跟随 strategy_weights.json 激活集合)
|
||||
|
||||
背景(2026-08-12 事项四:健康监控注册 p_oversold):
|
||||
evolution/health_monitor.py 能算策略健康度(实盘近7天浮盈 vs 回测基线5y 偏差 + 写 strategy_health 表 + 告警),
|
||||
@@ -31,6 +31,29 @@ def _singleton_guard(tag="health_monitor_daily.py"):
|
||||
sys.exit(0)
|
||||
|
||||
|
||||
def _load_active_strategies():
|
||||
"""从 strategy_weights.json 读当前激活策略(A股 active + 港股 markets.hk.active)。
|
||||
|
||||
2026-08-15 数据驱动改造:原硬编码 (v_weak, v_oversold),温区切换后激活策略变化
|
||||
但监控集合不跟随,导致激活策略无健康度数据。改为读路由输出,兜底旧二元组。
|
||||
"""
|
||||
import json as _json
|
||||
try:
|
||||
d = _json.loads(Path("/home/hmo/MoFin/data/strategy_weights.json").read_text(encoding="utf-8"))
|
||||
active = list(d.get("active") or [])
|
||||
active += list(((d.get("markets") or {}).get("hk") or {}).get("active") or [])
|
||||
seen, out = set(), []
|
||||
for v in active:
|
||||
if v and v not in seen:
|
||||
seen.add(v)
|
||||
out.append(v)
|
||||
if out:
|
||||
return out
|
||||
except Exception as e:
|
||||
print(f" [warn] 读 strategy_weights.json 失败({e}),兜底 v_weak/v_oversold", flush=True)
|
||||
return ["v_weak", "v_oversold"]
|
||||
|
||||
|
||||
def main():
|
||||
_fd = _singleton_guard()
|
||||
print(f"[health_monitor_daily] {datetime.now().strftime('%H:%M:%S')} 策略健康度监控开始", flush=True)
|
||||
@@ -44,9 +67,11 @@ def main():
|
||||
_spec.loader.exec_module(_m)
|
||||
run_health_check = _m.run_health_check
|
||||
|
||||
# 当前策略:v_weak(实盘在跑)+ v_oversold(新策略上线)
|
||||
# 数据驱动:监控集合 = 路由激活策略(A股+港股),随温区切换自动跟随
|
||||
versions = _load_active_strategies()
|
||||
print(f" 监控策略集合: {versions}", flush=True)
|
||||
results = {}
|
||||
for version in ("v_weak", "v_oversold"):
|
||||
for version in versions:
|
||||
try:
|
||||
r = run_health_check(strategy_version=version)
|
||||
results[version] = r
|
||||
|
||||
@@ -52,7 +52,9 @@ def get_trades_from_research(version):
|
||||
conn = sqlite3.connect(DB, timeout=30)
|
||||
conn.execute("PRAGMA busy_timeout=30000")
|
||||
rows = conn.execute(
|
||||
"SELECT results_json FROM strategy_research WHERE version=? ORDER BY period_tag DESC, created_at DESC LIMIT 1",
|
||||
# 2026-08-15:最长周期优先(LENGTH 降序:'10y' len3 > '5y'/'2y' len2)——
|
||||
# period_rollup 每日派生 1y/2y/5y 短窗行后,字典序 period_tag DESC 会错取短窗 trades
|
||||
"SELECT results_json FROM strategy_research WHERE version=? ORDER BY LENGTH(COALESCE(period_tag,'2y')) DESC, COALESCE(period_tag,'2y') DESC, created_at DESC LIMIT 1",
|
||||
(version,)
|
||||
).fetchall()
|
||||
conn.close()
|
||||
|
||||
@@ -0,0 +1,173 @@
|
||||
#!/usr/bin/env python3
|
||||
# -*- coding: utf-8 -*-
|
||||
"""strategy_period_rollup.py — 策略周期统计每日滚动派生(2026-08-15 问题2:策略数据每日自动更新)
|
||||
|
||||
背景:
|
||||
dashboard 研究 Tab 的 1y/2y/5y 列依赖 strategy_research 对应 period_tag 行,
|
||||
但 s2_panic/v_lurk_v3/v_mr4 等激活策略只有 10y 记录 → 列空、health_monitor 无 5y 基线。
|
||||
本脚本每日从「最长周期主记录」的 trades 按窗口切出 1y/2y/5y,派生行 upsert 回表。
|
||||
|
||||
原则:
|
||||
1. 窗口锚定主记录数据末端(max entry_date),不锚今天——回测未重跑时窗口不缩水
|
||||
2. 不覆盖真实研究记录:只 upsert 带 derived_from 标记的派生行;真实行存在则跳过
|
||||
3. 派生行含 calc_summary + portfolio_sim(10槽) + portfolio_full(100槽),
|
||||
口径与 strategy_lab.list_strategies 的 read-time 切窗一致(1m/6m/1y 已有同类逻辑)
|
||||
调度:每日盘后(策略路由后、regime_perf 前),hermes cron。
|
||||
"""
|
||||
import sys
|
||||
import json
|
||||
import sqlite3
|
||||
from pathlib import Path
|
||||
from datetime import datetime, timedelta
|
||||
from collections import Counter
|
||||
|
||||
_SCRIPT_DIR = Path(__file__).resolve().parent
|
||||
sys.path.insert(0, str(_SCRIPT_DIR))
|
||||
sys.path.insert(0, "/home/hmo/MoFin")
|
||||
|
||||
DB = "/home/hmo/MoFin/data/mofin.db"
|
||||
|
||||
# 目标派生周期(10y 作为主记录源,不派生)
|
||||
TARGET_PERIODS = {"1y": 365, "2y": 730, "5y": 1826}
|
||||
# 周期强度排序(选主记录用):越长越优先
|
||||
_PERIOD_RANK = {"10y": 6, "7y": 5, "5y": 4, "2y": 3, "1y": 2, "6m": 1, "1m": 0}
|
||||
|
||||
|
||||
def _period_rank(tag):
|
||||
return _PERIOD_RANK.get(tag or "2y", 3)
|
||||
|
||||
|
||||
def load_master_rows(conn):
|
||||
"""每个 (version, market) 取周期最长、含 trades 的记录作为主记录"""
|
||||
rows = conn.execute(
|
||||
"SELECT id, version, name, summary, hypothesis, parent, config_json, results_json, "
|
||||
"period, COALESCE(market,'all'), COALESCE(period_tag,'2y') FROM strategy_research"
|
||||
).fetchall()
|
||||
best = {}
|
||||
for r in rows:
|
||||
key = (r[1], r[9])
|
||||
rank = _period_rank(r[10])
|
||||
if key not in best or (rank, r[0]) > (best[key][0], best[key][1]):
|
||||
best[key] = (rank, r[0], r)
|
||||
return [v[2] for v in best.values()]
|
||||
|
||||
|
||||
def existing_real_periods(conn, version, market):
|
||||
"""该 (version, market) 已存在的【真实】period_tag(results_json 无 derived_from 标记)"""
|
||||
rows = conn.execute(
|
||||
"SELECT COALESCE(period_tag,'2y'), results_json FROM strategy_research "
|
||||
"WHERE version=? AND COALESCE(market,'all')=?",
|
||||
(version, market)).fetchall()
|
||||
real, derived = set(), set()
|
||||
for tag, rj in rows:
|
||||
try:
|
||||
d = json.loads(rj or "{}")
|
||||
except Exception:
|
||||
d = {}
|
||||
if d.get("derived_from"):
|
||||
derived.add(tag)
|
||||
else:
|
||||
real.add(tag)
|
||||
return real, derived
|
||||
|
||||
|
||||
def derive_period(master, target_tag, days, calc_summary, portfolio_sim):
|
||||
"""从主记录切窗派生 target_tag 记录。返回 (results_dict, period_str) 或 None"""
|
||||
res = json.loads(master["results_json"])
|
||||
trades = res.get("trades") or []
|
||||
if not trades:
|
||||
return None
|
||||
max_date = max(t.get("entry_date", "") for t in trades)
|
||||
if not max_date:
|
||||
return None
|
||||
cutoff = (datetime.strptime(max_date, "%Y-%m-%d") - timedelta(days=days)).strftime("%Y-%m-%d")
|
||||
sliced = [t for t in trades if t.get("entry_date", "") >= cutoff]
|
||||
if not sliced:
|
||||
return None
|
||||
for t in sliced:
|
||||
t.setdefault("boost", 1.0)
|
||||
capital = res.get("capital") or 1000000
|
||||
summary = calc_summary(sliced, capital)
|
||||
try:
|
||||
summary["portfolio"] = portfolio_sim(sliced, 1000000, max_positions=10)
|
||||
summary["portfolio_full"] = portfolio_sim(sliced, 1000000, max_positions=100)
|
||||
except Exception as e:
|
||||
print(f" [warn] portfolio_sim 失败({e})", flush=True)
|
||||
year_dist = dict(Counter(t["entry_date"][:4] for t in sliced if t.get("entry_date")))
|
||||
out = {
|
||||
"strategy": res.get("strategy") or master["version"],
|
||||
"strategy_name": res.get("strategy_name") or master["name"],
|
||||
"market": res.get("market") or master["market"],
|
||||
"period": f"{cutoff} ~ {max_date}",
|
||||
"period_tag": target_tag,
|
||||
"capital": capital,
|
||||
"total_stocks_screened": None, # 窗口派生无法还原,置空而非伪造
|
||||
"scored_events": None,
|
||||
"trades": sliced,
|
||||
"summary": summary,
|
||||
"year_dist": year_dist,
|
||||
"derived_from": master["period_tag"],
|
||||
"derived_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
|
||||
}
|
||||
return out
|
||||
|
||||
|
||||
def main():
|
||||
from strategy_lab import calc_summary, portfolio_sim
|
||||
|
||||
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
conn = sqlite3.connect(DB, timeout=60)
|
||||
conn.execute("PRAGMA busy_timeout=60000")
|
||||
|
||||
masters = load_master_rows(conn)
|
||||
print(f"[period_rollup] {now} 主记录 {len(masters)} 个 (version, market)", flush=True)
|
||||
|
||||
inserted = skipped_real = skipped_empty = 0
|
||||
for m in masters:
|
||||
master = {
|
||||
"id": m[0], "version": m[1], "name": m[2], "summary": m[3],
|
||||
"hypothesis": m[4], "parent": m[5], "config_json": m[6],
|
||||
"results_json": m[7], "period": m[8], "market": m[9], "period_tag": m[10],
|
||||
}
|
||||
real, derived = existing_real_periods(conn, master["version"], master["market"])
|
||||
for tag, days in TARGET_PERIODS.items():
|
||||
if tag == master["period_tag"]:
|
||||
continue # 主记录本身就是该周期
|
||||
if _period_rank(tag) >= _period_rank(master["period_tag"]):
|
||||
continue # 只向更短周期派生
|
||||
if tag in real:
|
||||
skipped_real += 1
|
||||
continue # 真实研究记录存在,不覆盖
|
||||
try:
|
||||
res = derive_period(master, tag, days, calc_summary, portfolio_sim)
|
||||
except Exception as e:
|
||||
print(f" [err] {master['version']}/{master['market']}/{tag}: {e}", flush=True)
|
||||
continue
|
||||
if not res:
|
||||
skipped_empty += 1
|
||||
continue
|
||||
# 删旧派生行,插新派生行(保持每日最新)
|
||||
conn.execute(
|
||||
"DELETE FROM strategy_research WHERE version=? AND COALESCE(market,'all')=? "
|
||||
"AND COALESCE(period_tag,'2y')=? AND results_json LIKE '%\"derived_from\"%'",
|
||||
(master["version"], master["market"], tag))
|
||||
conn.execute(
|
||||
"INSERT INTO strategy_research "
|
||||
"(version, name, summary, hypothesis, parent, config_json, results_json, "
|
||||
" analysis_json, period, created_at, market, period_tag) "
|
||||
"VALUES (?,?,?,?,?,?,?,?,?,?,?,?)",
|
||||
(master["version"], master["name"], master["summary"], master["hypothesis"],
|
||||
master["parent"], master["config_json"], json.dumps(res, ensure_ascii=False),
|
||||
json.dumps({"derived": True, "from": master["period_tag"]}, ensure_ascii=False),
|
||||
res["period"], now, master["market"], tag))
|
||||
inserted += 1
|
||||
s = res["summary"]
|
||||
print(f" + {master['version']:<14} {master['market']:<3} {tag}: "
|
||||
f"{s.get('total_trades')}笔 胜率{s.get('win_rate')}% (派生自{master['period_tag']})", flush=True)
|
||||
conn.commit()
|
||||
conn.close()
|
||||
print(f"[period_rollup] 完成: 派生写入 {inserted} 行 | 真实记录跳过 {skipped_real} | 空窗跳过 {skipped_empty}", flush=True)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user