实盘补齐大盘市场阶段判断(market_regime): 与回测_load_index_ctx同算法

- 新增 market_regime.py: sh000001 的 above_ma20/ma20_slope/roc/adx 计算+regime分类(trend_up/choppy/trend_down)+market_regime表
- market_watch.py: 每30分钟调度时同步更新 market_regime
- strategy_lifecycle.py: enrich_timing_signal 新增市场阶段因子+降级逻辑(震荡/下跌市买入信号降级为关注), regenerate_all/reassess_with_context 加载regime
This commit is contained in:
hmo
2026-08-02 16:29:45 +08:00
parent b676d85c41
commit 21401f91cf
3 changed files with 325 additions and 2 deletions
+259
View File
@@ -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_strengthbacktest_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 的 sh000001import_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())
+15
View File
@@ -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推送
+51 -2
View File
@@ -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