Files
MoFin/deploy/profile-scripts/anomaly_monitor.py
T

172 lines
6.4 KiB
Python

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""anomaly_monitor.py — 持仓股异动监控(2026-08-19 老莫:中科电气放量拉升无通知的教训)
检测持仓股/自选股的异动(价格急变/量能异常/突破),异动时 ACTION 级推送(不被信息淹没)。
异动判定(确定性代码,不依赖 LLM):
1. 单日涨跌幅 ≥ 阈值(默认 ±5%)→ 急涨/急跌
2. 量能异常:当日成交量 ≥ 5日均量 × 3 → 放量
3. 价格突破:创 20 日新高(或跌破 20 日新低)
4. 深套股反弹:浮亏 ≤ -20% 的持仓单日涨幅 ≥ 4%(深套异动,老莫关注)
推送:ACTION 级(alert_helper 直通),只对持仓股 + 自选股。
用法:cron 每 5 分钟跑(交易日盘中)
python3 anomaly_monitor.py
"""
import sqlite3, os, sys
from datetime import datetime
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent))
DB = "/home/hmo/MoFin/data/mofin.db"
THRESHOLD_CHG = 5.0 # 单日涨跌幅阈值 %
THRESHOLD_VOL = 3.0 # 量能倍数(vs 5日均量)
THRESHOLD_20D = True # 20日新高/新低检测
THRESHOLD_DEEP_REBOUND = 4.0 # 深套股反弹阈值 %
STATE_FILE = "/home/hmo/MoFin/data/anomaly_state.json"
def get_conn():
conn = sqlite3.connect(DB, timeout=30)
conn.execute("PRAGMA busy_timeout=30000")
return conn
def load_state():
try:
import json
with open(STATE_FILE, encoding="utf-8") as f:
return json.load(f)
except Exception:
return {"notified": {}} # {code_date: ts}
def save_state(st):
import json
try:
# 只保留最近 500 条
if len(st.get("notified", {})) > 500:
st["notified"] = dict(list(st["notified"].items())[-500:])
with open(STATE_FILE, "w", encoding="utf-8") as f:
json.dump(st, f)
except Exception:
pass
def push_anomaly(body):
"""ACTION 级推送(直通,不限速)"""
try:
from alert_helper import notify, ACTION
notify("持仓异动", body, level=ACTION)
return True
except Exception as e:
print(f" ⚠️ 异动推送失败: {e}", flush=True)
return False
def get_holdings():
"""持仓 + 自选股列表"""
conn = get_conn()
# 持仓(holdings 表)
holds = conn.execute("SELECT code, name, shares, cost, price FROM holdings WHERE is_active=1 AND shares>0").fetchall()
# 自选(holding_strategies 自选策略)
wl = conn.execute("SELECT code, name FROM holding_strategies WHERE status='active' AND decision_type='自选策略'").fetchall()
conn.close()
codes = set(h[0] for h in holds) | set(w[0] for w in wl)
names = {h[0]: h[1] or h[0] for h in holds}
names.update({w[0]: w[1] or w[0] for w in wl})
costs = {h[0]: h[2] for h in holds}
return codes, names, costs
def check_anomalies():
conn = get_conn()
codes, names, costs = get_holdings()
st = load_state()
notified = st.get("notified", {})
now = datetime.now().strftime("%Y-%m-%d %H:%M")
today = now[:10]
anomalies = []
for code in sorted(codes):
# 最新日线 + 昨收 + 前5日均量
row = conn.execute("""
SELECT date, close, volume FROM stock_daily
WHERE code=? ORDER BY date DESC LIMIT 6
""", (code,)).fetchall()
if len(row) < 2:
continue
latest = row[0] # (date, close, volume)
prev = row[1]
if not latest[1] or not prev[1] or prev[1] <= 0:
continue
# 当日涨跌幅
chg = (latest[1] - prev[1]) / prev[1] * 100
# 5日均量(排除当日)
vols = [r[2] for r in row[1:6] if r[2]]
avg_vol = sum(vols) / len(vols) if vols else 0
vol_ratio = latest[2] / avg_vol if avg_vol > 0 else 0
# 20日新高/新低
high20 = conn.execute("""
SELECT MAX(high), MIN(low) FROM (
SELECT high, low FROM stock_daily WHERE code=? ORDER BY date DESC LIMIT 20
)
""", (code,)).fetchone()
is_new_high = bool(high20 and high20[0] and latest[1] >= high20[0])
is_new_low = bool(high20 and high20[1] and latest[1] <= high20[1])
# 深套反弹(成本 vs 现价)
cost = costs.get(code) or 0
pnl = (latest[1] - cost) / cost * 100 if cost > 0 else 0
is_deep = pnl <= -20
# 判定异动
reasons = []
if abs(chg) >= THRESHOLD_CHG:
reasons.append(f"{'急涨' if chg > 0 else '急跌'}{abs(chg):.1f}%")
if vol_ratio >= THRESHOLD_VOL:
reasons.append(f"放量{vol_ratio:.1f}倍")
if THRESHOLD_20D and is_new_high:
reasons.append("创20日新高")
if THRESHOLD_20D and is_new_low:
reasons.append("破20日新低")
if is_deep and chg >= THRESHOLD_DEEP_REBOUND:
reasons.append(f"深套反弹{chg:.1f}%")
if reasons and latest[0] == today:
# 去重:同 code 同原因类型当天只推一次
key = f"{code}_{'|'.join(sorted(reasons))}_{today}"
if key not in notified or (datetime.now().timestamp() - notified[key] > 3600):
name = names.get(code, code)
cost_s = f" 成本{cost:.2f} 浮盈{pnl:+.1f}%" if cost > 0 else ""
body = (f"⚡ {name}({code}) 异动\n"
f"现价 {latest[1]:.2f} ({chg:+.1f}%) | 量能{vol_ratio:.1f}\n"
f"触发: {'、'.join(reasons)}{cost_s}")
anomalies.append(body)
notified[key] = datetime.now().timestamp()
print(f" ⚡ {code}: {'、'.join(reasons)}", flush=True)
save_state({"notified": notified})
conn.close()
if anomalies:
full = f"🔔 持仓/自选异动 {len(anomalies)}\n" + "\n\n".join(anomalies)
push_anomaly(full)
print(f"[异动] 推送 {len(anomalies)} 条", flush=True)
else:
print(f"[异动] {now} 无异动", flush=True)
return len(anomalies)
if __name__ == "__main__":
# 单例守卫(防 cron 重复)
import fcntl
lock_f = open("/tmp/anomaly_monitor.lock", "w")
try:
fcntl.flock(lock_f, fcntl.LOCK_EX | fcntl.LOCK_NB)
except Exception:
print("已有实例在运行,跳过", flush=True)
sys.exit(0)
try:
raise SystemExit(check_anomalies())
finally:
try:
fcntl.flock(lock_f, fcntl.LOCK_UN)
except Exception:
pass