diff --git a/deploy/profile-scripts/market_regime.py b/deploy/profile-scripts/market_regime.py new file mode 100644 index 00000000..05e5d02c --- /dev/null +++ b/deploy/profile-scripts/market_regime.py @@ -0,0 +1,259 @@ +#!/usr/bin/env python3 +"""market_regime.py — 大盘市场阶段判断(实盘版) + +解决实盘"无法判断当前是什么市"的断层: +回测用 _load_index_ctx('sh000001') 算 above_ma20/ma20_slope/roc/adx 做市场过滤 +(v_next4 要求 above_ma20=True + adx>=20 才出手,v_mr 用 mkt_mode='any'), +但实盘 strategy_lifecycle 只有当日涨跌幅 mood(±1%/1.5%→仓位乘数 0.8~1.1), +没有 MA20+ADX 的市场阶段判断 —— 本模块把回测算法原样搬到实盘。 + +算法与 /home/hmo/MoFin/strategy_lab.py 的 _load_index_ctx 完全一致: + - above_ma20 : 收盘价 > MA20 + - ma20_slope : MA20 相对 5 日前 MA20 的变化率(%) + - roc : 10 日 Rate of Change(%) + - adx : 简易趋势强度 (DI+ - DI-)/(DI+ + DI-) 归一化,14 日窗口 + (calc_trend_strength,backtest_framework.py 87 行) + +regime 分类(与回测 v_next4/v_mr 分工对齐): + - trend_up : above_ma20=True 且 adx>=20 → 趋势市,v_next4 主战场(追涨有效) + - choppy : adx<20 → 震荡市,v_mr 主战场(超跌反弹) + - trend_down : above_ma20=False 且 adx>=20 → 下跌趋势,v_mr 主战场(深超跌) + +用法: + python3 market_regime.py # 计算并写 market_regime 表(供定时调度) + python3 market_regime.py --print # 只打印当前市场状态 + 作为库: from market_regime import compute_regime, load_market_regime + +写表: market_regime(date PK, above_ma20, ma20_slope, roc, adx, regime, close, created_at) +数据源: mofin.db stock_daily 的 sh000001(import_full_stocks 每日收盘后更新, + 盘中用最新可得日线,未收盘日不计入最终判断,adx 用真实历史)。 +""" +import sys +import json +import sqlite3 +from pathlib import Path +from datetime import datetime + +# ── 路径注入:可被 deploy/profile-scripts 下脚本直接 import ── +_SCRIPT_DIR = Path(__file__).resolve().parent +_MOFIN_ROOT = _SCRIPT_DIR.parent.parent # deploy/profile-scripts → MoFin +for _p in (str(_SCRIPT_DIR), str(_MOFIN_ROOT)): + if _p not in sys.path: + sys.path.insert(0, _p) + +DB_PATH = Path(_MOFIN_ROOT) / "data" / "mofin.db" +INDEX_CODE = "sh000001" # 上证指数(A股) +INDEX_CODE_HK = "hkHSI" # 恒生指数(港股) + +# ADX 阈值与回测 v_next4 的 mkt_adx_min=20 对齐 +ADX_TREND_MIN = 20.0 + + +def calc_trend_strength(highs, lows, closes, n=14): + """简易趋势强度(替代 ADX): (DI+ - DI-) / (DI+ + DI-) 归一化 + 与 backtest_framework.calc_trend_strength 完全同算法(内联,避免跨目录依赖)""" + tr = [] + for i in range(len(highs)): + if i == 0: + tr.append(highs[i] - lows[i]) + else: + tr.append(max(highs[i] - lows[i], + abs(highs[i] - closes[i - 1]), + abs(lows[i] - closes[i - 1]))) + up = [highs[i] - highs[i - 1] for i in range(1, len(highs))] + down = [lows[i - 1] - lows[i] for i in range(1, len(lows))] + di_plus_raw = [0.0] * len(up) + di_minus_raw = [0.0] * len(up) + for i in range(len(up)): + if up[i] > down[i] and up[i] > 0: + di_plus_raw[i] = up[i] + if down[i] > up[i] and down[i] > 0: + di_minus_raw[i] = down[i] + tr_val = tr[i + 1] if (i + 1) < len(tr) else tr[-1] + if tr_val > 0: + di_plus_raw[i] = di_plus_raw[i] / tr_val * 100 + di_minus_raw[i] = di_minus_raw[i] / tr_val * 100 + result = [] + for i in range(len(di_plus_raw)): + if i < n - 1: + result.append(None) + else: + avg_plus = sum(di_plus_raw[i - n + 1:i + 1]) / n + avg_minus = sum(di_minus_raw[i - n + 1:i + 1]) / n + if avg_plus + avg_minus > 0: + dx = abs(avg_plus - avg_minus) / (avg_plus + avg_minus) * 100 + else: + dx = 0 + result.append(dx) + return [None] * (len(highs) - len(result)) + result + + +def calc_ma(series, n): + result = [] + for i in range(len(series)): + if i < n - 1: + result.append(None) + else: + result.append(sum(series[i - n + 1:i + 1]) / n) + return result + + +def calc_roc(series, n=10): + result = [] + for i in range(len(series)): + if i < n: + result.append(None) + else: + result.append((series[i] - series[i - n]) / series[i - n] * 100 if series[i - n] != 0 else 0) + return result + + +def compute_regime(index_code=INDEX_CODE, db_path=None, lookback_days=120): + """计算大盘市场状态(最近一个完整交易日)。 + + 返回 dict: + {date, close, above_ma20, ma20_slope, roc, adx, regime, computed_at} + 数据不足(<30行)或取数失败时返回 None。 + """ + db = db_path or DB_PATH + conn = sqlite3.connect(str(db), timeout=5) + try: + rows = conn.execute( + "SELECT date, close, high, low FROM stock_daily " + "WHERE code=? ORDER BY date DESC LIMIT ?", + (index_code, lookback_days)).fetchall() + finally: + conn.close() + if not rows or len(rows) < 30: + return None + # 升序处理(与回测 prepare_bars 一致) + rows = list(reversed(rows)) + dates = [r[0] for r in rows] + closes = [r[1] for r in rows] + highs = [r[2] for r in rows] + lows = [r[3] for r in rows] + + ma20 = calc_ma(closes, 20) + trend = calc_trend_strength(highs, lows, closes) + roc = calc_roc(closes) + i = len(closes) - 1 # 最新一日 + close = closes[i] + m20 = ma20[i] + above_ma20 = (close > m20) if m20 else None + # ma20_slope: 相对 5 日前 MA20 的变化率 + slope = None + if i >= 5 and ma20[i - 5] and ma20[i - 5] > 0 and m20: + slope = round((m20 - ma20[i - 5]) / ma20[i - 5] * 100, 3) + adx = trend[i] if i < len(trend) else None + roc_v = roc[i] if i < len(roc) else None + + # regime 分类 + if above_ma20 is True and adx is not None and adx >= ADX_TREND_MIN: + regime = "trend_up" + elif adx is not None and adx < ADX_TREND_MIN: + regime = "choppy" + else: + regime = "trend_down" + + return { + "date": dates[i], + "close": close, + "above_ma20": above_ma20, + "ma20_slope": slope, + "roc": round(roc_v, 3) if roc_v is not None else None, + "adx": round(adx, 2) if adx is not None else None, + "regime": regime, + "computed_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), + } + + +def save_regime(regime, db_path=None): + """写入 market_regime 表(按 date 去重,同日覆盖)""" + if not regime: + return False + db = db_path or DB_PATH + conn = sqlite3.connect(str(db), timeout=5) + try: + conn.execute(""" + CREATE TABLE IF NOT EXISTS market_regime ( + date TEXT PRIMARY KEY, + above_ma20 INTEGER, + ma20_slope REAL, + roc REAL, + adx REAL, + regime TEXT, + close REAL, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP + ) + """) + conn.execute(""" + INSERT OR REPLACE INTO market_regime + (date, above_ma20, ma20_slope, roc, adx, regime, close, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP) + """, ( + regime["date"], + 1 if regime["above_ma20"] else 0, + regime.get("ma20_slope"), + regime.get("roc"), + regime.get("adx"), + regime["regime"], + regime.get("close"), + )) + conn.commit() + finally: + conn.close() + return True + + +def load_market_regime(db_path=None): + """读取最新 market_regime(供 strategy_lifecycle 等消费方调用)""" + db = db_path or DB_PATH + conn = sqlite3.connect(str(db), timeout=5) + try: + row = conn.execute( + "SELECT date, above_ma20, ma20_slope, roc, adx, regime, close " + "FROM market_regime ORDER BY date DESC LIMIT 1").fetchone() + finally: + conn.close() + if not row: + return None + return { + "date": row[0], + "above_ma20": bool(row[1]), + "ma20_slope": row[2], + "roc": row[3], + "adx": row[4], + "regime": row[5], + "close": row[6], + } + + +REGIME_DESC = { + "trend_up": "趋势市(大盘MA20上方+ADX≥20):v_next4 追涨主战场", + "choppy": "震荡市(ADX<20):v_mr 超跌反弹主战场", + "trend_down": "下跌趋势(大盘MA20下方+ADX≥20):v_mr 深超跌主战场", +} + + +def main(): + regime = compute_regime() + if "--print" in sys.argv or not regime: + if not regime: + print("sh000001 数据不足,无法计算市场状态") + return 1 + print(json.dumps(regime, ensure_ascii=False, indent=2)) + print(f"判断: {REGIME_DESC.get(regime['regime'], '')}") + return 0 + ok = save_regime(regime) + if ok: + print(f"[market_regime] {regime['date']} → {regime['regime']} " + f"(above_ma20={regime['above_ma20']} adx={regime['adx']} " + f"slope={regime['ma20_slope']} roc={regime['roc']})") + else: + print("[market_regime] 写入失败") + return 1 + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/deploy/profile-scripts/market_watch.py b/deploy/profile-scripts/market_watch.py index e6c60d67..a51c31de 100644 --- a/deploy/profile-scripts/market_watch.py +++ b/deploy/profile-scripts/market_watch.py @@ -214,6 +214,21 @@ def main(): print(f"[DB] 写入失败(JSON 不受影响): {msg}", flush=True) conn.close() + # ── 大盘市场阶段(market_regime):与回测 _load_index_ctx 同算法 ── + # 回测用 above_ma20 + adx 做市场过滤(v_next4 要求趋势市才出手), + # 实盘此前只有当日涨跌幅 mood,缺 MA20+ADX 的市场阶段判断 —— 此处补齐。 + try: + import market_regime as _mr + _regime = _mr.compute_regime() + if _regime: + _mr.save_regime(_regime) + print(f"[regime] {_regime['date']} → {_regime['regime']} " + f"(above_ma20={_regime['above_ma20']} adx={_regime['adx']})", flush=True) + else: + print("[regime] sh000001 数据不足,跳过", flush=True) + except Exception as _e: + print(f"[regime] 计算失败: {_e}", flush=True) + # 靜默:只寫文件,不輸出到stdout,避免cron推送 diff --git a/deploy/profile-scripts/strategy_lifecycle.py b/deploy/profile-scripts/strategy_lifecycle.py index 9834177a..15ac37dd 100644 --- a/deploy/profile-scripts/strategy_lifecycle.py +++ b/deploy/profile-scripts/strategy_lifecycle.py @@ -1834,8 +1834,9 @@ def enrich_timing_signal(base_signal, macro_desc="", sector_note="", fundamentals=None, news_sentiment=None, timing_signal_override=None, portfolio_context=None, - rr_ratio=0): # 2026-06-24 新参:盈亏比约束 - """多因子合成timing_signal——大盘+行业+基本面+技术+组合风险+盈亏比 + rr_ratio=0, + market_regime=None): # 2026-08-02 新参:大盘市场阶段(trend_up/choppy/trend_down) + """多因子合成timing_signal——大盘+行业+基本面+技术+组合风险+盈亏比+市场阶段 返回 (enriched_signal, factors_list) - enriched_signal: 可读的多因子信号描述 @@ -1857,6 +1858,20 @@ def enrich_timing_signal(base_signal, macro_desc="", sector_note="", elif macro_desc and macro_desc != "宏观未加载": factors.append("大盘中性") + # 1.5 市场阶段因子(2026-08-02 新增——实盘补齐回测的 market_regime) + # 回测用 sh000001 的 above_ma20+adx 过滤市场(v_next4 仅趋势市出手,v_mr 超跌市主力), + # 实盘此前只有当日涨跌幅 mood,缺中期市场阶段判断 —— 此处用 market_regime 表补齐。 + if market_regime: + _rg = market_regime.get("regime", "") + if _rg == "trend_up": + factors.append("趋势市") + elif _rg == "trend_down": + factors.append("下跌趋势市") + elif _rg == "choppy": + factors.append("震荡市") + else: + factors.append(f"市场:{_rg}") + # 2. 行业因子 if sector_note: # 把"行业X大跌3%+"简化为"行业偏弱","行业X大涨3%+"简化为"行业偏强" @@ -1974,6 +1989,18 @@ def enrich_timing_signal(base_signal, macro_desc="", sector_note="", clean_signal = "信号不充分" factors.append("RR过低降级") + # 6.5 市场阶段降级(2026-08-02 新增——实盘复现回测 v_next4 的市场过滤) + # 回测实证:v_next4 全市场10年103笔买入 100% 发生在大盘 MA20 上方(趋势市), + # 震荡市/下跌趋势市追涨胜率低。实盘补齐该过滤: + # trend_up → 正常,买入/加仓信号放行(v_next4 主战场) + # choppy → 买入/加仓降级为"关注"(v_mr 超跌反弹主战场,趋势追涨让位) + # trend_down → 买入/加仓降级为"关注"(v_mr 深超跌主战场,禁止追涨) + if clean_signal in buy_signals and market_regime: + _rg = market_regime.get("regime", "") + if _rg in ("choppy", "trend_down"): + clean_signal = "关注" + factors.append(f"{'震荡市' if _rg=='choppy' else '下跌市'}买入降级") + return clean_signal, factors @@ -2005,6 +2032,14 @@ def reassess_with_context(code, name, price, cost, shares, current_action, news_sentiment = {} fund = {} + # 大盘市场阶段(market_regime)— 与回测 _load_index_ctx 同算法,趋势市放行追涨 + market_regime = None + try: + import market_regime as _mr + market_regime = _mr.load_market_regime() + except Exception: + pass # market_regime 不可用时不阻塞单只重评 + # ── DSA 集成:注入大盘复盘 + 新闻情报 ────────────────────────── try: from mo_bridge import enrich_analysis_context @@ -2027,6 +2062,7 @@ def reassess_with_context(code, name, price, cost, shares, current_action, news_sentiment=news_sentiment, portfolio_context=_get_portfolio_risk_state(), rr_ratio=result.get("rr_ratio", 0), + market_regime=market_regime, ) result["timing_signal"] = enriched result["signal_factors"] = factors @@ -2234,6 +2270,18 @@ def regenerate_all(stdout=True): sectors_found = sum(1 for c in all_stocks if stock_sector_map.get(c)) print(f" 市场参考: {market_mood} 上涨比{market_breadth}% 已匹配{sectors_found}/{total}只个股行业") + # 加载大盘市场阶段(market_regime)— 与回测 _load_index_ctx 同算法, + # 趋势市放行追涨买入,震荡/下跌市降级买入信号(v_next4/v_mr 分工的实盘落地) + market_regime = None + try: + import market_regime as _mr + market_regime = _mr.load_market_regime() + if stdout and market_regime: + print(f" 市场阶段: {market_regime.get('date')} " + f"{_mr.REGIME_DESC.get(market_regime.get('regime'), market_regime.get('regime'))}") + except Exception: + pass # market_regime 不可用时不阻塞主流程 + # 批量预取所有价格(一次API调用 vs 之前N次) prices_map = batch_fetch_prices(list(all_stocks.keys())) if stdout: @@ -2372,6 +2420,7 @@ def regenerate_all(stdout=True): fundamentals=fund, news_sentiment=news_sentiment, rr_ratio=result.get("rr_ratio", 0), + market_regime=market_regime, ) result["timing_signal"] = enriched