feat: 策略执行调度器 strategy_executor——按当前温区激活策略动态调度对应选股scanner(v_weak→mr_scanner/v_oversold→predictive_oversold/s2_panic→s2_scanner/v_next4→leader_scanner), 替代分散独立scanner cron, 温区休眠策略不扫描
This commit is contained in:
@@ -0,0 +1,149 @@
|
||||
#!/usr/bin/env python3
|
||||
# -*- coding: utf-8 -*-
|
||||
"""strategy_executor.py — 策略执行调度器(2026-08-13 温区自适应核心)
|
||||
|
||||
核心(老莫要求):根据当前温区激活策略,调度对应的选股/买卖/重评任务。
|
||||
|
||||
温区自适应数据流:
|
||||
regime_tracker(K=5 平滑温区) + temp_band(温度) + regime_perf(实测)
|
||||
→ strategy_router(当前温区激活策略 strategy_weights.json)
|
||||
→ strategy_executor(按激活策略调度对应 scanner/监控/重评)
|
||||
|
||||
调度逻辑:
|
||||
- 读 strategy_weights.json 的 active(当前温区激活策略集合)
|
||||
- 对每个激活策略,调度其对应的选股扫描器(scanner)
|
||||
- 非激活策略的 scanner 不跑(温区休眠)
|
||||
|
||||
策略→scanner 映射:
|
||||
v_weak → mr_scanner.py (choppy/trend_down 超跌确认)
|
||||
v_oversold → predictive_oversold_scanner.py(trend_down 预测超跌反弹)
|
||||
s2_panic → s2_scanner.py (trend_down 恐慌买超跌)
|
||||
v_mr/v_mr_sel → mr_scanner.py (复用 v_weak 扫描器,同家族)
|
||||
v_next4/v8.1 → leader_scanner.py (trend_up 龙头识别)
|
||||
|
||||
用法:cron 每30分钟盘中跑(替代各 scanner 分散 cron)
|
||||
python3 strategy_executor.py
|
||||
"""
|
||||
import sys
|
||||
import json
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
from datetime import datetime
|
||||
|
||||
_SCRIPT_DIR = Path(__file__).resolve().parent
|
||||
sys.path.insert(0, str(_SCRIPT_DIR))
|
||||
sys.path.insert(0, "/home/hmo/MoFin")
|
||||
|
||||
WEIGHTS_FILE = "/home/hmo/MoFin/data/strategy_weights.json"
|
||||
LOG = "/home/hmo/MoFin/gateway/logs/strategy_executor.log"
|
||||
|
||||
# 策略 → 选股扫描器映射
|
||||
STRATEGY_SCANNER = {
|
||||
"v_weak": "mr_scanner.py",
|
||||
"v_oversold": "predictive_oversold_scanner.py",
|
||||
"p_oversold": "predictive_oversold_scanner.py",
|
||||
"s2_panic": "s2_scanner.py",
|
||||
"v_mr": "mr_scanner.py",
|
||||
"v_mr_sel": "mr_scanner.py",
|
||||
"v_mr2": "mr_scanner.py",
|
||||
"v_mr3": "mr_scanner.py",
|
||||
"v_mr4": "mr_scanner.py",
|
||||
"v_lurk_v1": "mr_scanner.py",
|
||||
"v_lurk_v2": "mr_scanner.py",
|
||||
"v_lurk_v3": "mr_scanner.py",
|
||||
"v_next4": "leader_scanner.py",
|
||||
"v8.1": "leader_scanner.py",
|
||||
"v_next5": "leader_scanner.py",
|
||||
"v_combo": "leader_scanner.py",
|
||||
}
|
||||
|
||||
|
||||
def log(msg):
|
||||
line = f"[{datetime.now().isoformat(timespec='seconds')}] {msg}"
|
||||
print(line, flush=True)
|
||||
try:
|
||||
Path(LOG).parent.mkdir(parents=True, exist_ok=True)
|
||||
with open(LOG, "a", encoding="utf-8") as f:
|
||||
f.write(line + "\n")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def load_active():
|
||||
"""读当前温区激活策略"""
|
||||
try:
|
||||
p = Path(WEIGHTS_FILE)
|
||||
if p.exists():
|
||||
d = json.loads(p.read_text(encoding="utf-8"))
|
||||
return d.get("active") or [], d.get("state", "unknown")
|
||||
except Exception as e:
|
||||
log(f"加载 weights 失败: {e}")
|
||||
return [], "unknown"
|
||||
|
||||
|
||||
def is_trading_time():
|
||||
"""盘中时段检查(9:30-11:30 / 13:00-15:00 连续竞价)"""
|
||||
now = datetime.now()
|
||||
t = now.strftime("%H:%M")
|
||||
wd = now.weekday()
|
||||
if wd >= 5:
|
||||
return False
|
||||
return ("09:30" <= t <= "11:30") or ("13:00" <= t <= "15:00")
|
||||
|
||||
|
||||
def run_scanner(scanner_name):
|
||||
"""执行一个扫描器(带超时防挂)"""
|
||||
script = _SCRIPT_DIR / scanner_name
|
||||
if not script.exists():
|
||||
log(f" ⚠️ {scanner_name} 不存在,跳过")
|
||||
return False
|
||||
log(f" ▶ 执行 {scanner_name}")
|
||||
try:
|
||||
r = subprocess.run(
|
||||
[sys.executable, str(script)],
|
||||
capture_output=True, text=True, timeout=600,
|
||||
cwd=str(_SCRIPT_DIR),
|
||||
)
|
||||
out = (r.stdout or "").strip()[:300]
|
||||
err = (r.stderr or "").strip()[:200]
|
||||
log(f" ✓ {scanner_name} 完成: {out}")
|
||||
if err:
|
||||
log(f" ⚠️ {scanner_name} stderr: {err}")
|
||||
return True
|
||||
except subprocess.TimeoutExpired:
|
||||
log(f" ⏱ {scanner_name} 超时(600s)跳过")
|
||||
return False
|
||||
except Exception as e:
|
||||
log(f" ✗ {scanner_name} 异常: {e}")
|
||||
return False
|
||||
|
||||
|
||||
def main():
|
||||
if not is_trading_time():
|
||||
print(f"[SKIP] {datetime.now().strftime('%H:%M')} 非连续竞价时段", flush=True)
|
||||
return 0
|
||||
|
||||
active, regime = load_active()
|
||||
log(f"温区: {regime} | 激活策略: {active}")
|
||||
|
||||
# 去重:同一 scanner 只跑一次(多策略映射同一扫描器)
|
||||
scanners = set()
|
||||
for strat in active:
|
||||
sc = STRATEGY_SCANNER.get(strat)
|
||||
if sc:
|
||||
scanners.add(sc)
|
||||
log(f"调度扫描器: {sorted(scanners)}")
|
||||
|
||||
if not scanners:
|
||||
log(" 无激活策略的扫描器(当前温区无对应选股任务)")
|
||||
return 0
|
||||
|
||||
for sc in sorted(scanners):
|
||||
run_scanner(sc)
|
||||
|
||||
log("调度完成")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
Reference in New Issue
Block a user