feat(messenger): 全部77个活跃cron脚本接入统一消息通道(install_stdio_hook)

This commit is contained in:
xxm
2026-08-21 22:34:25 +08:00
parent 5e0141eb28
commit 20c05b6e7c
77 changed files with 2083 additions and 1544 deletions
+220 -213
View File
@@ -1,213 +1,220 @@
#!/usr/bin/env python3
"""capital_flow_collector.py — 个股资金流数据采集器
每30分钟拉一次持仓+自选的超大单/大单/中单/小单资金流向。
输出到 capital_flow_cache.json 供 price_monitor 和报告使用。
API: push2his.eastmoney.com 个股资金流日线
"""
import json, os, sys, time, urllib.request
from datetime import datetime
from urllib.request import urlopen, Request
from concurrent.futures import ThreadPoolExecutor, as_completed
from threading import Semaphore
from mofin_db import get_conn, write_capital_flow_cache
DATA_DIR = "/home/hmo/web-dashboard/data"
CACHE_PATH = f"{DATA_DIR}/capital_flow_cache.json"
UA = "Mozilla/5.0"
# 限速器:最多5个并发,每请求后强制间隔0.3s
RATE_LIMIT = Semaphore(5)
MIN_INTERVAL = 0.3
_last_req = 0
def _rate_limited_request(url, referer="https://data.eastmoney.com/"):
"""带速率限制的HTTP GET,用Semaphore控制并发数"""
global _last_req
with RATE_LIMIT:
elapsed = time.time() - _last_req
if elapsed < MIN_INTERVAL:
time.sleep(MIN_INTERVAL - elapsed)
proxy_handler = urllib.request.ProxyHandler({})
opener = urllib.request.build_opener(proxy_handler)
req = Request(url, headers={"User-Agent": UA, "Referer": referer})
try:
resp = opener.open(req, timeout=8)
_last_req = time.time()
return json.loads(resp.read().decode("utf-8"))
except Exception:
return None
# eastmoney secid: 1=上海 0=深圳
def secid(code):
code = str(code).strip()
if code.startswith(("6", "9")):
return f"1.{code}"
return f"0.{code}"
def fetch_flow(code, days=5):
"""拉取个股近N日资金流(Sina MoneyFlow ssl_qsfx_zjlrqs,最新在前 —
eastmoney 在 246 不可达、ssl_qsfx_lscjfb 数据停在2020年(2026-07-22 实证)"""
code = str(code).strip()
if code.startswith(("6", "9")):
dm = f"sh{code}"
elif len(code) == 5 and code[0] in "01":
return None # 港股 sina 不支持资金流
else:
dm = f"sz{code}"
url = ("https://vip.stock.finance.sina.com.cn/quotes_service/api/json_v2.php/"
f"MoneyFlow.ssl_qsfx_zjlrqs?daima={dm}")
data = _rate_limited_request(url, referer="https://finance.sina.com.cn")
if not data or not isinstance(data, list) or not data:
return None
result = []
for d in data[:days]: # 最新在前,取近N日
try:
r0n = float(d.get("r0_net", 0) or 0)
r1n = float(d.get("r1_net", 0) or 0)
r2n = float(d.get("r2_net", 0) or 0)
r3n = float(d.get("r3_net", 0) or 0)
result.append({
"date": d.get("opendate", ""),
"main_net": r0n + r1n,
"super_large": r0n,
"large": r1n,
"medium": r2n,
"small": r3n,
})
except Exception:
continue
result.reverse() # 转为时间升序(analyze_flow 以 [-1] 为最新日)
return result or None
def fetch_flow_intraday(code):
"""拉取当日分时资金流(用于盘中判断)"""
sid = secid(code)
url = f"http://push2.eastmoney.com/api/qt/stock/fflow/kline/get?secid={sid}&fields1=f1,f2,f3,f7&fields2=f51,f52,f53,f54,f55,f56,f57&klt=1&lmt=120"
try:
resp = urlopen(url, timeout=5)
data = json.loads(resp.read().decode("utf-8"))
klines = data.get("data", {}).get("klines", [])
if not klines:
return None
latest = klines[-1].split(",")
return {
"main_net": float(latest[1]),
"super_large": float(latest[2]),
"large": float(latest[3]),
}
except:
return None
def analyze_flow(flow_data):
"""分析资金流模式"""
if not flow_data or len(flow_data) < 2:
return {}
result = {"alerts": [], "pattern": ""}
# 最近两日对比
d1 = flow_data[-1] # 最新日
d2 = flow_data[-2] # 前一日
# 超大单信号
sl1 = d1["super_large"]
sl2 = d2["super_large"]
# 连续形态判断
main_trend = sum(d["main_net"] for d in flow_data[-3:])
sl_trend = sum(d["super_large"] for d in flow_data[-3:])
# 1. 主力连续流入
if main_trend > 50000000 and sl1 > 0 and sl2 > 0:
result["pattern"] = "主力持续流入"
result["alerts"].append("主力连续3日净流入")
# 2. 超大单突然转向(连续流入→流出 或 流出→流入)
if sl1 * sl2 < 0: # 方向反转
if sl1 > 0 and sl2 < 0:
result["pattern"] = "超大单由出转入"
result["alerts"].append("超大单转为净买入(暗示消息即将落地)")
elif sl1 < 0 and sl2 > 0:
result["pattern"] = "超大单由入转出"
result["alerts"].append("超大单转为净卖出(利好出货嫌疑)")
# 3. 价格与资金流背离(缺当前价格作比较,在主脚本中完成)
# 4. 单日暴量
max_sl = max(abs(d["super_large"]) for d in flow_data)
if max_sl == abs(sl1) and abs(sl1) > 100000000:
result["pattern"] = "单日资金暴量"
result["alerts"].append(f"今日超大单异常: {sl1/100000000:.2f}亿")
return result
def main():
codes = set()
# 读取持仓+自选(从DB直接读,替代已删除的mo_data)
try:
import sqlite3
_db = sqlite3.connect("/home/hmo/MoFin/data/mofin.db")
for row in _db.execute("SELECT DISTINCT code FROM holdings WHERE is_active=1").fetchall():
if row[0]: codes.add(row[0])
for row in _db.execute("SELECT DISTINCT code FROM holding_strategies WHERE status='active' AND decision_type='自选策略'").fetchall():
if row[0]: codes.add(row[0])
_db.close()
except:
pass
all_flows = {}
# 并行抓取:ThreadPoolExecutor + 内置限速器(Semaphore 5 + 0.3s间隔)
code_list = sorted(codes)
if not code_list:
print("[capital_flow] 无代码需要采集")
return
def fetch_one(code):
flow = fetch_flow(code, days=5)
if flow:
analysis = analyze_flow(flow) or {}
# 聚合成 build_prompt 需要的字段(此前只有 alerts/pattern12维资金面恒为0
net = round(sum(d["main_net"] for d in flow) / 1e4, 1)
main_f = round(sum(d["super_large"] for d in flow) / 1e4, 1)
retail = round(sum(d["medium"] + d["small"] for d in flow) / 1e4, 1)
analysis["net_flow"] = net
analysis["main_force"] = main_f
analysis["retail_flow"] = retail
analysis["trend"] = analysis.get("pattern") or \
("主力流入" if main_f > 0 else ("主力流出" if main_f < 0 else "中性"))
return (code, {
"updated_at": datetime.now().strftime("%Y-%m-%d %H:%M"),
"flow": flow,
"analysis": analysis,
})
return (code, None)
with ThreadPoolExecutor(max_workers=5) as pool:
futures = {pool.submit(fetch_one, c): c for c in code_list}
for f in as_completed(futures):
code, result = f.result()
if result:
all_flows[code] = result
# 写缓存
cache = {
"updated_at": datetime.now().strftime("%Y-%m-%d %H:%M"),
"stocks": all_flows,
}
# 写 DB(替代 capital_flow_cache.json
conn = get_conn()
write_capital_flow_cache(conn, cache)
conn.close()
print(f"[capital_flow] {len(all_flows)}/{len(code_list)}只更新完成")
# 2026-08-13 资金流突变处理(老莫设计):正面突变→选股分析,负面突变(持仓)→重评+XMPP报告
try:
from fund_flow_alert import process_flow_alerts
process_flow_alerts(all_flows)
except Exception as e:
print(f"[资金流突变处理异常] {e}", flush=True)
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""capital_flow_collector.py — 个股资金流数据采集器
每30分钟拉一次持仓+自选的超大单/大单/中单/小单资金流向。
输出到 capital_flow_cache.json 供 price_monitor 和报告使用。
API: push2his.eastmoney.com 个股资金流日线
"""
import json, os, sys, time, urllib.request
from datetime import datetime
from urllib.request import urlopen, Request
from concurrent.futures import ThreadPoolExecutor, as_completed
from threading import Semaphore
from mofin_db import get_conn, write_capital_flow_cache
# ── 消息通道统一路由(broadcast/xmpp by delivery) ──
try:
from messenger import install_stdio_hook as _msh
_msh()
except Exception:
pass
DATA_DIR = "/home/hmo/web-dashboard/data"
CACHE_PATH = f"{DATA_DIR}/capital_flow_cache.json"
UA = "Mozilla/5.0"
# 限速器:最多5个并发,每请求后强制间隔0.3s
RATE_LIMIT = Semaphore(5)
MIN_INTERVAL = 0.3
_last_req = 0
def _rate_limited_request(url, referer="https://data.eastmoney.com/"):
"""带速率限制的HTTP GET,用Semaphore控制并发数"""
global _last_req
with RATE_LIMIT:
elapsed = time.time() - _last_req
if elapsed < MIN_INTERVAL:
time.sleep(MIN_INTERVAL - elapsed)
proxy_handler = urllib.request.ProxyHandler({})
opener = urllib.request.build_opener(proxy_handler)
req = Request(url, headers={"User-Agent": UA, "Referer": referer})
try:
resp = opener.open(req, timeout=8)
_last_req = time.time()
return json.loads(resp.read().decode("utf-8"))
except Exception:
return None
# eastmoney secid: 1=上海 0=深圳
def secid(code):
code = str(code).strip()
if code.startswith(("6", "9")):
return f"1.{code}"
return f"0.{code}"
def fetch_flow(code, days=5):
"""拉取个股近N日资金流(Sina MoneyFlow ssl_qsfx_zjlrqs,最新在前 —
eastmoney 在 246 不可达、ssl_qsfx_lscjfb 数据停在2020年(2026-07-22 实证)"""
code = str(code).strip()
if code.startswith(("6", "9")):
dm = f"sh{code}"
elif len(code) == 5 and code[0] in "01":
return None # 港股 sina 不支持资金流
else:
dm = f"sz{code}"
url = ("https://vip.stock.finance.sina.com.cn/quotes_service/api/json_v2.php/"
f"MoneyFlow.ssl_qsfx_zjlrqs?daima={dm}")
data = _rate_limited_request(url, referer="https://finance.sina.com.cn")
if not data or not isinstance(data, list) or not data:
return None
result = []
for d in data[:days]: # 最新在前,取近N日
try:
r0n = float(d.get("r0_net", 0) or 0)
r1n = float(d.get("r1_net", 0) or 0)
r2n = float(d.get("r2_net", 0) or 0)
r3n = float(d.get("r3_net", 0) or 0)
result.append({
"date": d.get("opendate", ""),
"main_net": r0n + r1n,
"super_large": r0n,
"large": r1n,
"medium": r2n,
"small": r3n,
})
except Exception:
continue
result.reverse() # 转为时间升序(analyze_flow 以 [-1] 为最新日)
return result or None
def fetch_flow_intraday(code):
"""拉取当日分时资金流(用于盘中判断)"""
sid = secid(code)
url = f"http://push2.eastmoney.com/api/qt/stock/fflow/kline/get?secid={sid}&fields1=f1,f2,f3,f7&fields2=f51,f52,f53,f54,f55,f56,f57&klt=1&lmt=120"
try:
resp = urlopen(url, timeout=5)
data = json.loads(resp.read().decode("utf-8"))
klines = data.get("data", {}).get("klines", [])
if not klines:
return None
latest = klines[-1].split(",")
return {
"main_net": float(latest[1]),
"super_large": float(latest[2]),
"large": float(latest[3]),
}
except:
return None
def analyze_flow(flow_data):
"""分析资金流模式"""
if not flow_data or len(flow_data) < 2:
return {}
result = {"alerts": [], "pattern": ""}
# 最近两日对比
d1 = flow_data[-1] # 最新日
d2 = flow_data[-2] # 前一日
# 超大单信号
sl1 = d1["super_large"]
sl2 = d2["super_large"]
# 连续形态判断
main_trend = sum(d["main_net"] for d in flow_data[-3:])
sl_trend = sum(d["super_large"] for d in flow_data[-3:])
# 1. 主力连续流入
if main_trend > 50000000 and sl1 > 0 and sl2 > 0:
result["pattern"] = "主力持续流入"
result["alerts"].append("主力连续3日净流入")
# 2. 超大单突然转向(连续流入→流出 或 流出→流入)
if sl1 * sl2 < 0: # 方向反转
if sl1 > 0 and sl2 < 0:
result["pattern"] = "超大单由出转入"
result["alerts"].append("超大单转为净买入(暗示消息即将落地)")
elif sl1 < 0 and sl2 > 0:
result["pattern"] = "超大单由入转出"
result["alerts"].append("超大单转为净卖出(利好出货嫌疑)")
# 3. 价格与资金流背离(缺当前价格作比较,在主脚本中完成)
# 4. 单日暴量
max_sl = max(abs(d["super_large"]) for d in flow_data)
if max_sl == abs(sl1) and abs(sl1) > 100000000:
result["pattern"] = "单日资金暴量"
result["alerts"].append(f"今日超大单异常: {sl1/100000000:.2f}亿")
return result
def main():
codes = set()
# 读取持仓+自选(从DB直接读,替代已删除的mo_data)
try:
import sqlite3
_db = sqlite3.connect("/home/hmo/MoFin/data/mofin.db")
for row in _db.execute("SELECT DISTINCT code FROM holdings WHERE is_active=1").fetchall():
if row[0]: codes.add(row[0])
for row in _db.execute("SELECT DISTINCT code FROM holding_strategies WHERE status='active' AND decision_type='自选策略'").fetchall():
if row[0]: codes.add(row[0])
_db.close()
except:
pass
all_flows = {}
# 并行抓取:ThreadPoolExecutor + 内置限速器(Semaphore 5 + 0.3s间隔)
code_list = sorted(codes)
if not code_list:
print("[capital_flow] 无代码需要采集")
return
def fetch_one(code):
flow = fetch_flow(code, days=5)
if flow:
analysis = analyze_flow(flow) or {}
# 聚合成 build_prompt 需要的字段(此前只有 alerts/pattern12维资金面恒为0
net = round(sum(d["main_net"] for d in flow) / 1e4, 1)
main_f = round(sum(d["super_large"] for d in flow) / 1e4, 1)
retail = round(sum(d["medium"] + d["small"] for d in flow) / 1e4, 1)
analysis["net_flow"] = net
analysis["main_force"] = main_f
analysis["retail_flow"] = retail
analysis["trend"] = analysis.get("pattern") or \
("主力流入" if main_f > 0 else ("主力流出" if main_f < 0 else "中性"))
return (code, {
"updated_at": datetime.now().strftime("%Y-%m-%d %H:%M"),
"flow": flow,
"analysis": analysis,
})
return (code, None)
with ThreadPoolExecutor(max_workers=5) as pool:
futures = {pool.submit(fetch_one, c): c for c in code_list}
for f in as_completed(futures):
code, result = f.result()
if result:
all_flows[code] = result
# 写缓存
cache = {
"updated_at": datetime.now().strftime("%Y-%m-%d %H:%M"),
"stocks": all_flows,
}
# 写 DB(替代 capital_flow_cache.json
conn = get_conn()
write_capital_flow_cache(conn, cache)
conn.close()
print(f"[capital_flow] {len(all_flows)}/{len(code_list)}只更新完成")
# 2026-08-13 资金流突变处理(老莫设计):正面突变→选股分析,负面突变(持仓)→重评+XMPP报告
try:
from fund_flow_alert import process_flow_alerts
process_flow_alerts(all_flows)
except Exception as e:
print(f"[资金流突变处理异常] {e}", flush=True)
if __name__ == "__main__":
main()