From 2af2e934085dc1e71aac768e891e7f00fd33d4db Mon Sep 17 00:00:00 2001 From: hmo Date: Sun, 2 Aug 2026 20:27:31 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20v=5Fmr=20=E5=AE=9E=E7=9B=98=E6=89=AB?= =?UTF-8?q?=E6=8F=8F=E5=99=A8=E2=80=94=E2=80=94=E5=9B=9E=E6=B5=8B=E5=9B=A0?= =?UTF-8?q?=E5=AD=90=E5=8E=9F=E6=A0=B7=E8=90=BD=E5=9C=B0,=E8=85=BE?= =?UTF-8?q?=E8=AE=AFqfq=E6=95=B0=E6=8D=AE=E6=BA=90=E4=B8=8Estock=5Fdaily?= =?UTF-8?q?=E9=9B=B6=E5=81=8F=E5=B7=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - mr_scanner.py: v_mr 6条件筛选(bias60<-10%/RSI<42/ret60<-15%/mom20<5%/额<1500万/rsi_delta>=2) - 股票池=stock_daily distinct code,与回测run_mr_backtest完全同口径(含ST,不含退市/301新创业板) - 数据源腾讯qfq日K: 因子与回测零偏差(实测bias60/ret60/mom20 Δ=0.000000) - regime门控: trend_down/choppy才扫; 幂等保护: 当天已有候选则跳过 - market_watch.py: 接入mr_scanner调用(每30分钟调度,幂等防重复) - 首轮实盘验证: 2026-07-31 trend_down市命中3只(*ST四通/*ST冀凯/*ST椰岛),全为回测验证过的ST股 --- deploy/profile-scripts/market_watch.py | 21 ++ deploy/profile-scripts/mr_scanner.py | 365 +++++++++++++++++++++++++ 2 files changed, 386 insertions(+) create mode 100644 deploy/profile-scripts/mr_scanner.py diff --git a/deploy/profile-scripts/market_watch.py b/deploy/profile-scripts/market_watch.py index a51c31de..74dd77b7 100644 --- a/deploy/profile-scripts/market_watch.py +++ b/deploy/profile-scripts/market_watch.py @@ -229,6 +229,27 @@ def main(): except Exception as _e: print(f"[regime] 计算失败: {_e}", flush=True) + # ── v_mr 均值回复扫描(每天最多一次,regime 门控内置)── + # trend_down/choppy 时扫描深超跌小盘票,写 candidates(sector='v_mr') + # mr_scanner 自带幂等保护:当天已有 v_mr 候选则跳过,避免重复全扫描 + try: + import subprocess as _sp + import sys as _sys + _script = Path(__file__).parent / "mr_scanner.py" + if _script.exists(): + _r = _sp.run( + [_sys.executable, str(_script), "--top", "15"], + capture_output=True, text=True, timeout=1500) + if _r.stdout: + for _line in _r.stdout.strip().splitlines(): + print(f"[mr_scanner] {_line}", flush=True) + if _r.returncode != 0 and _r.stderr: + print(f"[mr_scanner] stderr: {_r.stderr[:500]}", flush=True) + else: + print("[mr_scanner] 脚本不存在,跳过", flush=True) + except Exception as _e: + print(f"[mr_scanner] 调用失败: {_e}", flush=True) + # 靜默:只寫文件,不輸出到stdout,避免cron推送 diff --git a/deploy/profile-scripts/mr_scanner.py b/deploy/profile-scripts/mr_scanner.py new file mode 100644 index 00000000..f10a5447 --- /dev/null +++ b/deploy/profile-scripts/mr_scanner.py @@ -0,0 +1,365 @@ +#!/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()