v1.6.0: 目标跟踪(概念/主题/股票) - watch_targets表+概念智能体跟踪(关键词检索+受益个股+影响度判定+邮件通知), 自动化页加持仓管理和目标管理, 报告支持股票/概念双类型

This commit is contained in:
2026-08-19 23:49:37 +08:00
parent b34006c0d8
commit 9406f9de43
7 changed files with 389 additions and 65 deletions
+166 -25
View File
@@ -261,8 +261,8 @@ def track_stock(code, focus=""):
"ratings": seg["ratings"], "holdings": seg["holdings"],
}
execute(
"INSERT INTO tracking_reports(code, stock_name, industry, report, meta, sources, status, created_at) "
"VALUES(?,?,?,?,?,?,'done',datetime('now','localtime'))",
"INSERT INTO tracking_reports(target_type, code, stock_name, industry, report, meta, sources, status, created_at) "
"VALUES('stock',?,?,?,?,?,?,'done',datetime('now','localtime'))",
(code, s["name"], s["industry"], clean_report, json.dumps(meta, ensure_ascii=False),
json.dumps(sources, ensure_ascii=False)))
rid = query_one("SELECT MAX(id) id FROM tracking_reports")["id"]
@@ -287,60 +287,201 @@ def _first_heading(text):
def _notify_if_significant(rid, stock, meta):
"""影响度达阈值且开启通知 → 邮件"""
"""影响度达阈值且开启通知 → 邮件(股票/概念通用)"""
cfg = tracking_config()
try:
if cfg["notify"] and int(meta["impact_score"]) >= int(cfg["impact_threshold"]):
if cfg["notify"] and int(meta.get("impact_score", 0)) >= int(cfg["impact_threshold"]):
from engine.notifier import send_email
mc = mail_config()
nc = meta.get("news_counts") or {}
detail = (f"个股资讯 {nc.get('direct', 0)} 条 / 上游 {nc.get('upstream', 0)} 条 / "
f"下游 {nc.get('downstream', 0)} 条 / 同业 {nc.get('peers', 0)}"
if "upstream" in nc else
f"相关资讯 {nc.get('direct', 0)} 条 / 语义检索 {nc.get('rag', 0)} 条 / "
f"受益个股 {nc.get('related', 0)}")
send_email(
f"[持仓跟踪] {stock['name']} 出现{meta.get('change_kind','')}动态(影响度{meta['impact_score']}",
f"[持仓跟踪] {stock['name']} 出现{meta.get('change_kind', '')}动态(影响度{meta.get('impact_score', 0)}",
f"""<html><body style="font-family:Microsoft YaHei;padding:20px;background:#f5f6f8;">
<div style="max-width:640px;margin:auto;background:#fff;border-radius:8px;border:1px solid #e5e7eb;overflow:hidden;">
<div style="background:#1e293b;color:#fff;padding:14px 20px;font-size:17px;font-weight:bold;">🧭 持仓跟踪 · {stock['name']}{stock['code']}</div>
<div style="background:#1e293b;color:#fff;padding:14px 20px;font-size:17px;font-weight:bold;">🧭 智能体跟踪 · {stock['name']}{stock['code']}</div>
<div style="padding:16px 20px;">
<p><b>影响度:</b>{meta['impact_score']}/100{'🔴 重大' if meta['impact_score']>=65 else '🟡 关注'}<br>
<b>性质:</b>{meta.get('change_kind','')} <b>显著性:</b>{meta.get('significance','')}</p>
<p style="font-size:15px;"><b>摘要:</b>{meta.get('summary','')}</p>
<p style="color:#555;"><b>产业链趋势:</b>{meta.get('chain_trend','')}</p>
<p style="color:#888;font-size:12px;">个股资讯 {meta['news_counts']['direct']} 条 / 上游 {meta['news_counts']['upstream']} 条 / 下游 {meta['news_counts']['downstream']} 条 / 同业 {meta['news_counts']['peers']}</p>
<p><b>影响度:</b>{meta.get('impact_score', 0)}/100{'🔴 重大' if int(meta.get('impact_score', 0)) >= 65 else '🟡 关注'}<br>
<b>性质:</b>{meta.get('change_kind', '')} <b>显著性:</b>{meta.get('significance', '')}</p>
<p style="font-size:15px;"><b>摘要:</b>{meta.get('summary', '')}</p>
<p style="color:#555;"><b>趋势判断</b>{meta.get('chain_trend', '')}</p>
<p style="color:#888;font-size:12px;">{detail}</p>
</div></div></body></html>""",
cfg=mc)
set_tracking_state(last_alert=int(meta["impact_score"]))
set_tracking_state(last_alert=int(meta.get("impact_score", 0)))
except Exception as e:
log.warning("track notify fail: %s", e)
# ===================================================================== 概念/主题跟踪
def track_concept(name, keywords="", focus=""):
"""跟踪一个概念/主题:采集相关新闻 → 受益个股梳理 → LLM 深度分析 → 影响度判定"""
kw = keywords.strip() or name
kw_list = [k.strip() for k in kw.replace("", ",").split(",") if k.strip()]
kws = " ".join(kw_list)
# 1) 采集:DB 关键词检索 + RAG 向量检索
conds, args = [], []
for k in kw_list[:5]:
conds.append("(title LIKE ? OR content LIKE ?)")
args += [f"%{k}%", f"%{k}%"]
args.append(15)
db_news = query(
f"SELECT id,title,content,source,category,publish_date,sentiment,related_stocks FROM news "
f"WHERE ({' OR '.join(conds)}) AND publish_date >= date('now','-45 day') "
f"ORDER BY publish_date DESC LIMIT ?", args)
rag_hits = _rag_news_text(f"{name} 概念主题 {kws} 最新动态 政策 催化", top_k=8)
db_text = _fmt_news(db_news)
rag_text = "\n".join(rag_hits) or " (暂无)"
# 2) 受益个股:从 DB 新闻 related_stocks + RAG 命中元数据 code 汇总
codes = set()
for n in db_news:
for c in (n["related_stocks"] or "").split(","):
if c:
codes.add(c)
try:
for h in query_vectors(f"{name} {kws} 受益个股", n_results=8, name=CHROMA_NEWS_COLLECTION):
c = (h.get("metadata") or {}).get("code")
if c:
codes.add(c)
except Exception:
pass
related = query("SELECT code, name, industry FROM stocks WHERE code IN (%s)" %
",".join(["?"] * len(codes)), tuple(codes)) if codes else []
related_txt = "".join(f"{r['name']}({r['code']},{r['industry']})" for r in related) or "暂无"
# 3) 提示词
prompt = f"""你是资深题材/概念跟踪分析师,正在深度跟踪【{name}】这一概念主题。
【概念/主题】{name}
【检索关键词】{kws}
【最新相关资讯】
{db_text}
{rag_text}
【关联受益个股】{related_txt}
【输出要求】
第一步,先输出一个 json 代码块(必须最先输出):
```json
{{"significance":"high|medium|low","impact_score":0到100的整数,"change_kind":"利好/利空/中性/震荡","summary":"一句话总结","chain_trend":"题材趋势判断"}}
```
第二步,输出 Markdown 分析报告,结构如下:
## 一、概念主题最新动态
## 二、核心驱动与催化(政策/产业/事件)
## 三、受益标的梳理(关联个股及逻辑)
## 四、市场情绪与资金动向
## 五、风险提示
## 关注要点(3-5条)
分析须严格基于资讯,避免编造。impact_score >=65 视为重大变化。"""
try:
reply = llm_chat([
{"role": "system", "content": "你是一名严谨专业的题材与概念跟踪分析师。"},
{"role": "user", "content": prompt},
]).strip()
if not reply:
raise RuntimeError("LLM 返回为空")
judge = parse_judge(reply)
clean_report = strip_json_block(reply)
summary = (judge or {}).get("summary") or _first_heading(clean_report)
meta = {
"significance": (judge or {}).get("significance", "medium"),
"impact_score": int((judge or {}).get("impact_score", 50)),
"change_kind": (judge or {}).get("change_kind", "中性"),
"summary": summary,
"chain_trend": (judge or {}).get("chain_trend", ""),
"news_counts": {"direct": len(db_news), "rag": len(rag_hits), "related": len(related)},
"keywords": kws,
}
sources = {"keywords": kws, "db_news": db_text, "rag_news": rag_text,
"related": related_txt}
cid = f"CONCEPT:{name}"
execute(
"INSERT INTO tracking_reports(target_type, code, stock_name, industry, report, meta, sources, status, created_at) "
"VALUES('concept',?,?,?,?,?,?,'done',datetime('now','localtime'))",
(cid, name, "概念/主题", clean_report, json.dumps(meta, ensure_ascii=False),
json.dumps(sources, ensure_ascii=False)))
rid = query_one("SELECT MAX(id) id FROM tracking_reports")["id"]
_notify_if_significant(rid, {"name": name, "code": cid}, meta)
return {"ok": True, "report_id": rid, "meta": meta}
except Exception as e:
log.exception("track concept %s fail", name)
return {"error": str(e)}
# ===================================================================== 目标管理
def list_targets():
return query("SELECT * FROM watch_targets ORDER BY type, id")
def add_target(ttype, name, code="", keywords=""):
name = name.strip()
if not name:
return {"error": "名称不能为空"}
execute("INSERT INTO watch_targets(type, code, name, keywords) VALUES(?,?,?,?)",
(ttype, code, name, keywords))
return {"ok": True}
def delete_target(tid):
execute("DELETE FROM watch_targets WHERE id=?", (tid,))
return {"ok": True}
# ===================================================================== 批量与调度
def track_watchlist(progress=None):
"""串行跟踪自选股(持仓)。同一时刻只允许一个跟踪任务(防重复)"""
def track_all(progress=None):
"""跟踪全部目标:持仓(自选股)+ 目标(概念/主题/股票)"""
global _cycle_running
if _cycle_running:
return {"tracked": 0, "msg": "已有跟踪任务进行中,请稍后再试"}
_cycle_running = True
try:
stocks = query("SELECT w.code, s.name FROM watchlist w JOIN stocks s ON s.code=w.code ORDER BY w.added_at")
if not stocks:
return {"tracked": 0, "msg": "自选股为空,请先在股票池添加"}
targets = []
# 持仓(自选股)
for w in query("SELECT w.code, s.name FROM watchlist w JOIN stocks s ON s.code=w.code ORDER BY w.added_at"):
targets.append(("stock", w["code"], w["name"]))
# 目标(股票/概念)
for t in query("SELECT id, type, code, name, keywords FROM watch_targets WHERE enabled=1 ORDER BY id"):
if t["type"] == "stock" and t["code"]:
targets.append(("stock", t["code"], t["name"]))
else:
targets.append(("concept", t["name"], t["keywords"] or t["name"]))
# 去重
seen, uniq = set(), []
for typ, key, name in targets:
u = (typ, key)
if u in seen:
continue
seen.add(u)
uniq.append((typ, key, name))
if not uniq:
return {"tracked": 0, "msg": "暂无跟踪目标:请添加持仓/自选股或概念主题目标"}
results = []
for i, st in enumerate(stocks):
r = track_stock(st["code"])
results.append({"code": st["code"], "name": st["name"], **r})
set_tracking_state(last_run=time.strftime("%Y-%m-%d %H:%M:%S"), last_stock=st["name"])
for i, (typ, key, name) in enumerate(uniq):
if typ == "stock":
r = track_stock(key)
else:
r = track_concept(key, name)
results.append({"type": typ, "name": name, **r})
set_tracking_state(last_run=time.strftime("%Y-%m-%d %H:%M:%S"), last_stock=name)
if progress:
progress(i + 1, len(stocks))
progress(i + 1, len(uniq))
return {"tracked": len(results), "results": results}
finally:
_cycle_running = False
def latest_reports(code, limit=5):
return query("SELECT id, code, stock_name, industry, meta, status, created_at "
return query("SELECT id, target_type, code, stock_name, industry, meta, status, created_at "
"FROM tracking_reports WHERE code=? ORDER BY id DESC LIMIT ?", (code, limit))
def list_reports(limit=30):
return query("SELECT id, code, stock_name, industry, meta, status, created_at "
return query("SELECT id, target_type, code, stock_name, industry, meta, status, created_at "
"FROM tracking_reports ORDER BY id DESC LIMIT ?", (limit,))
@@ -371,7 +512,7 @@ class TrackingThread(threading.Thread):
continue
if cfg["enabled"]:
try:
r = track_watchlist()
r = track_all()
log.info("tracking cycle: %s", r)
except Exception as e:
log.warning("tracking cycle error: %s", e)