#!/usr/bin/env python3 """mr_scanner.py — v_mr 均值回复策略实盘扫描器 回测侧 v_mr(strategy_lab.run_mr_backtest / mr_engine_v2.py)已 10y 全市场验证 (21690笔 / WR 53.6% / avg +3.34%),但实盘此前没有该策略的扫描机制 —— 本脚本把 v_mr 的入场筛选原样搬到实盘,按日扫描全市场深超跌小盘票。 与回测完全一致的入场条件(mr_cfg v2 定稿): 1. bias60 < -10% : 收盘价在 MA60 下方超 10%(单边,无下限) 2. RSI14 < 42 : 超卖 3. ret60 < -15% : 60 日跌幅超 15%(单边,无下限) 4. mom20 < 5% : 20 日低动量(还没启动反弹) 5. amount20 < 1500万 : 小盘(20 日均成交额,百万元) 6. rsi_delta >= 2 : RSI 5 日回升 ≥2(止跌回升确认) 出场建议(exit_cfg v2):tp=+18% / sl=-8% / max_hold=25 交易日 市场门控(2026-08-02 新增): - 只读 market_regime 表,regime 为 trend_down 或 choppy 时启用扫描 (v_mr 主战场:下跌趋势 + 震荡市;趋势市让位 v_next4) - trend_up 时跳过(v_next4 追涨主战场,不扫超跌) 数据源:Sina 240 分钟线(日K),datalen=120(覆盖 MA60 + 60日回看 + RSI 收敛) 指标算法:与 backtest_framework.py 完全一致(内联,零偏差) 输出:candidates 表(sector='v_mr'),与 accumulation_scanner 同 UPSERT 模式 用法: python3 mr_scanner.py # 完整扫描(regime 门控) python3 mr_scanner.py --force # 忽略 regime 门控强制扫描 python3 mr_scanner.py --top N # 输出前 N 只(默认 10) """ import sys, json, urllib.request, re, time, sqlite3 from pathlib import Path from datetime import datetime DB_PATH = Path("/home/hmo/MoFin/data/mofin.db") UA = "Mozilla/5.0" # ── v_mr 参数(与 mr_engine_v2.py 注册的 mr_cfg/exit_cfg 完全一致)── MR_CFG = { "bias_max": -10, # MA60 下方超 10%(单边无下限) "rsi_max": 42, # RSI 超卖 "ret_max": -15, # 60 日跌超 15%(单边无下限) "mom20_max": 5, # 20 日低动量 "amount_max": 15, # 百万元 = 1500 万日成交额 "rsi_delta_min": 2, # RSI 5 日回升 ≥2 "mkt_mode": "any", # 实盘由 regime 门控替代(trend_down/choppy 才扫) } EXIT_CFG = {"tp_pct": 0.18, "sl_pct": 0.08, "max_hold_days": 25} TOP_N = 10 # ── 技术指标(与 backtest_framework.py 完全同算法)── 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_rsi(series, n=14): deltas = [series[i] - series[i - 1] for i in range(1, len(series))] gains = [d if d > 0 else 0 for d in deltas] losses = [-d if d < 0 else 0 for d in deltas] result = [None] * (n + 1) avg_gain = sum(gains[:n]) / n avg_loss = sum(losses[:n]) / n if avg_loss == 0: result.append(100) else: rs = avg_gain / avg_loss result.append(100 - 100 / (1 + rs)) for i in range(n, len(gains)): avg_gain = (avg_gain * (n - 1) + gains[i]) / n avg_loss = (avg_loss * (n - 1) + losses[i]) / n if avg_loss == 0: result.append(100) else: rs = avg_gain / avg_loss result.append(100 - 100 / (1 + rs)) while len(result) < len(series): result.insert(0, None) return result[:len(series)] # ── 数据获取 ── def fetch_tx_klines(code, datalen=120): """腾讯前复权日K(qfq),与 stock_daily 数据零偏差,返回 [{date,open,close,high,low,volume}]""" raw = str(code).strip() if raw.startswith(("6", "9")): prefix = "sh" elif raw.startswith(("0", "3")): prefix = "sz" else: return None url = f"http://ifzq.gtimg.cn/appstock/app/fqkline/get?param={prefix}{raw},day,,,{datalen},qfq" try: req = urllib.request.Request(url, headers={"User-Agent": UA}) opener = urllib.request.build_opener(urllib.request.ProxyHandler({})) with opener.open(req, timeout=8) as r: text = r.read().decode("utf-8", errors="replace").strip() data = json.loads(text) node = data.get("data", {}).get(f"{prefix}{raw}", {}) bars = node.get("qfqday") or node.get("day") or [] if not bars or len(bars) < 70: return None result = [] for b in bars: if len(b) < 6: continue result.append({ "date": b[0][:10], "open": float(b[1]), "close": float(b[2]), "high": float(b[3]), "low": float(b[4]), "volume": float(b[5]), # 手 }) return result except Exception: return None # 兼容别名(供外部引用) fetch_sina_klines = fetch_tx_klines def get_stock_pool(): """待扫描股票池:stock_daily 的 distinct code(与回测 run_mr_backtest 完全同口径) 回测股票池 = SELECT DISTINCT sd.code FROM stock_daily(4266只,含300/688, 不含301新创业板——数据源未收录)。实盘扫描用同一口径,保证 v_mr 信号覆盖的股票都是回测验证过的。 """ conn = sqlite3.connect(str(DB_PATH), timeout=5) try: existing = set() for r in conn.execute("SELECT code FROM holding_strategies WHERE status='active'"): existing.add(str(r[0])) for r in conn.execute("SELECT code FROM holdings WHERE is_active=1"): existing.add(str(r[0])) # 与回测完全一致:stock_daily 有K线的股票(回测 universe='a' 排除5位港股) all_stocks = [str(r[0]) for r in conn.execute("SELECT DISTINCT code FROM stock_daily").fetchall()] finally: conn.close() # 只留 A 股(6位数字),排除港股(5位0开头)—— 与回测 is_hk_code 逻辑一致 a_stocks = [c for c in all_stocks if len(c) == 6 and c.isdigit()] return a_stocks, existing def load_regime(): """读取 market_regime 最新状态""" try: conn = sqlite3.connect(str(DB_PATH), timeout=5) row = conn.execute( "SELECT date, above_ma20, adx, regime FROM market_regime " "ORDER BY date DESC LIMIT 1").fetchone() conn.close() if row: return {"date": row[0], "above_ma20": bool(row[1]), "adx": row[2], "regime": row[3]} except Exception: pass return None # ── v_mr 筛选(与回测 run_mr_backtest 的 6 条件一致)── def check_vmr(klines): """对单只股票做 v_mr 入场筛选。命中返回信号 dict,否则 None。 klines 为升序日K(Sina 返回顺序)。""" if not klines or len(klines) < 70: return None closes = [k["close"] for k in klines] highs = [k["high"] for k in klines] lows = [k["low"] for k in klines] i = len(klines) - 1 # 最新一日 close = closes[i] if close <= 0: return None ma60 = calc_ma(closes, 60) rsi_all = calc_rsi(closes) rsi = rsi_all[i] if i < len(rsi_all) else None m60 = ma60[i] if not m60 or m60 <= 0 or rsi is None: return None # 1. MA60 下方超跌(单边) bias60 = (close - m60) / m60 * 100 if bias60 > MR_CFG["bias_max"]: # 只拦上沿(-10%),无下限 return None # 2. RSI 超卖 if rsi > MR_CFG["rsi_max"]: return None # 3. 60 日跌幅(单边) prev60 = closes[i - 60] if i >= 60 else 0 prev_ret60 = (close - prev60) / prev60 * 100 if prev60 > 0 else 0 if prev_ret60 > MR_CFG["ret_max"]: # 只拦上沿(-15%),无下限 return None # 4. 20 日低动量 prev20 = closes[i - 20] if i >= 20 else 0 mom20 = (close - prev20) / prev20 * 100 if prev20 > 0 else 0 if mom20 > MR_CFG["mom20_max"]: return None # 5. 小盘(20 日均成交额,百万元) # 腾讯日K volume 单位=手(×100股),成交额=手×100×均价 amt20 = [] for k in klines[max(0, i - 19):i + 1]: avg_px = (k["high"] + k["low"] + k["close"]) / 3 amt20.append(k["volume"] * 100 * avg_px) # 元 amt_valid = [a for a in amt20 if a and a > 0] amount_ma20 = sum(amt_valid) / len(amt_valid) if amt_valid else 0 amount_ma20_m = amount_ma20 / 1e6 # 元 → 百万元 if MR_CFG["amount_max"] is not None and amount_ma20_m > MR_CFG["amount_max"]: return None # 6. RSI 5 日回升(止跌确认) rsi0 = rsi_all[i - 5] if i >= 5 else None rsi_delta = (rsi - rsi0) if rsi0 is not None else 0 if rsi_delta < MR_CFG["rsi_delta_min"]: return None # 命中 → 出场建议(与回测 exit_cfg 一致) tp_pct = EXIT_CFG["tp_pct"] sl_pct = EXIT_CFG["sl_pct"] target = round(close * (1 + tp_pct), 2) stop = round(close * (1 - sl_pct), 2) return { "price": close, "bias60": round(bias60, 2), "rsi": round(rsi, 2), "prev_ret60": round(prev_ret60, 2), "mom20": round(mom20, 2), "amount_ma20": round(amount_ma20_m, 2), "rsi_delta": round(rsi_delta, 2), "target": target, "stop_loss": stop, "date": klines[i]["date"], } def main(): force = "--force" in sys.argv top_n = TOP_N if "--top" in sys.argv: try: top_n = int(sys.argv[sys.argv.index("--top") + 1]) except (ValueError, IndexError): pass print(f"[MR] {datetime.now().strftime('%H:%M')} v_mr 实盘扫描开始", flush=True) # ── regime 门控 ── regime = load_regime() if regime: rg = regime["regime"] print(f" 市场阶段: {regime['date']} → {rg} (adx={regime['adx']})", flush=True) if rg not in ("trend_down", "choppy") and not force: print(f" ⏭ {rg} 非 v_mr 主战场(trend_down/choppy 才扫),跳过", flush=True) return if rg not in ("trend_down", "choppy"): print(f" ⚠ --force 强制扫描(当前 {rg})", flush=True) else: print(" ⚠ market_regime 不可用,默认执行扫描(v_mr 全市场可用)", flush=True) # ── 幂等检查:当天已扫过 v_mr 则跳过(避免 30 分钟调度重复全扫描)── import sqlite3 as _sq _conn = _sq.connect(str(DB_PATH), timeout=5) try: _today = datetime.now().strftime("%Y-%m-%d") _n = _conn.execute( "SELECT COUNT(*) FROM candidates WHERE sector='v_mr' AND substr(created_at,1,10)=?", (_today,)).fetchone()[0] except Exception: _n = 0 _conn.close() if _n > 0 and not force: print(f" 已有 {_n} 条今日 v_mr 候选,跳过重复扫描(--force 可强制)", flush=True) return # ── 股票池 ── all_stocks, existing = get_stock_pool() print(f" 股票池: {len(all_stocks)}只A股, 已有策略: {len(existing)}只", flush=True) if not all_stocks: print(" ⚠ stocks 表为空", flush=True) return # 并发拉日K(ThreadPool 8 并发,Sina 单只逐个拉) from concurrent.futures import ThreadPoolExecutor, as_completed pool = [c for c in all_stocks if c not in existing] found = [] done = 0 with ThreadPoolExecutor(max_workers=8) as ex: fut_map = {ex.submit(fetch_sina_klines, c): c for c in pool} for fut in as_completed(fut_map): code = fut_map[fut] done += 1 klines = fut.result() if klines: sig = check_vmr(klines) if sig: found.append((code, sig)) if done % 400 == 0: print(f" 已扫描 {done}/{len(pool)}", flush=True) print(f" 命中 v_mr 条件: {len(found)} 只", flush=True) # 排序:超跌越深越优先(bias60 越小越靠前) found.sort(key=lambda x: x[1]["bias60"]) # ── 写 candidates 表(UPSERT,保留计算列)── conn = sqlite3.connect(str(DB_PATH), timeout=5) inserted = 0 for code, sig in found[:top_n]: name = code try: r = conn.execute("SELECT name FROM stocks WHERE code=?", (code,)).fetchone() if r and r[0]: name = r[0] except Exception: pass price = sig["price"] entry_low = round(price * 0.97, 2) entry_high = round(price * 1.02, 2) sl = sig["stop_loss"] tp = sig["target"] reasons = (f"v_mr超跌(bias60={sig['bias60']}% rsi={sig['rsi']} " f"ret60={sig['prev_ret60']}% mom20={sig['mom20']}% " f"额{sig['amount_ma20']}M rsi_delta={sig['rsi_delta']})") # 检查是否已在 candidates 且未 promoted exists = conn.execute( "SELECT code FROM candidates WHERE code=? AND (promoted IS NULL OR promoted=0)", (code,)).fetchone() if exists: continue conn.execute( "INSERT INTO candidates (code, name, sector, reason, " "entry_range, stop_loss, target, created_at) " "VALUES (?,?,?,?,?,?,?,datetime('now','localtime')) " "ON CONFLICT(code) DO UPDATE SET " "name=excluded.name, sector=excluded.sector, reason=excluded.reason, " "entry_range=excluded.entry_range, stop_loss=excluded.stop_loss, target=excluded.target", (code, name, "v_mr", reasons, f"{entry_low}~{entry_high}", sl, tp)) inserted += 1 print(f" 🟢 {code} {name} 价{price} bias60={sig['bias60']}% " f"rsi={sig['rsi']} ret60={sig['prev_ret60']}% {reasons}", flush=True) conn.commit() conn.close() print(f" ✅ 新增 {inserted} 只 v_mr 候选(前 {top_n})", flush=True) if __name__ == "__main__": main()