diff --git a/src/analysis/event_impact.py b/src/analysis/event_impact.py new file mode 100644 index 0000000..45b0460 --- /dev/null +++ b/src/analysis/event_impact.py @@ -0,0 +1,271 @@ +# -*- coding: utf-8 -*- +""" +事件影响力分析:新闻/快讯 → 市场影响推断 → 个股推荐卡。 +LLM 只负责解读文本与推断方向;候选股必须落在真实行情快照内(grounding), +购入区间基于真实现价计算,预计收益标注为推测。 +""" +import json +import sqlite3 +import traceback +from datetime import datetime, timedelta + +import requests + +from src.analysis.llm_client import LlmClient, LlmError + +SYSTEM_PROMPT = """你是 A 股事件影响分析师。系统提供:若干条最新财经快讯/新闻、当日主力资金流入TOP10、\ +今日龙虎榜摘要、上证指数状态,以及全市场个股现价快照(部分)。 +任务:判断这些消息中是否包含值得关注的**事件**;若有,推断该事件对哪个品类(板块)的股票\ +产生什么方向的影响,并从"现价快照"给出的股票中选出最可能受益的标的。 + +铁律: +1. 只允许从系统提供的现价快照中选股票,禁止编造代码/名称。 +2. 证据必须引用所给新闻原文片段。 +3. 没有值得分析的事件时输出 has_event=false。 +4. 输出严格 JSON,无 markdown 代码块: +{"has_event": true|false, + "event_summary": "≤80字事件摘要", + "direction": "利好|利空|中性", + "sectors": ["受影响板块", 最多3个], + "confidence": 0-100, + "reasoning": "≤150字影响传导逻辑", + "stocks": [{"name":"股票名","code":"6位代码","reason":"≤60字推荐理由", + "stars": 1到5的整数, + "expected_return_pct": 预计5日收益百分数(可为负), + "buy_zone_low": 建议购入区间下限, "buy_zone_high": 建议购入区间上限}] + stocks 最多 3 只;buy_zone 必须参考该股现价(快照中有 price 字段)。""" + +ANALYZE_COOLDOWN_S = 180 +MAX_NEWS_PER_RUN = 8 + + +class EventImpactService: + + def __init__(self, emit, db_path): + self.emit = emit + self.db_path = str(db_path) + self.llm = LlmClient() + self.last_analyzed_at = None # 只分析该时间之后的快讯 + self.last_run = 0.0 + + def _conn(self): + return sqlite3.connect(self.db_path, check_same_thread=False) + + def _ensure_tables(self, conn): + conn.execute("""CREATE TABLE IF NOT EXISTS event_impact ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ts TEXT NOT NULL, + event_summary TEXT NOT NULL, + direction TEXT, + sectors TEXT, + confidence INTEGER, + reasoning TEXT, + stocks TEXT, + evidence TEXT + )""") + conn.execute("""CREATE TABLE IF NOT EXISTS news_analyzed ( + event_time TEXT NOT NULL, + content TEXT NOT NULL, + analyzed_at TEXT NOT NULL, + PRIMARY KEY (event_time, content) + )""") + conn.commit() + + def has_new_news(self, conn) -> bool: + cutoff = (datetime.now() - timedelta(hours=3)).strftime('%Y-%m-%d %H:%M:%S') + n = conn.execute( + "SELECT COUNT(*) FROM news_flash WHERE event_time >= ? AND event_time > COALESCE(" + " (SELECT MAX(analyzed_at) FROM news_analyzed), '1970-01-01')", (cutoff,)).fetchone()[0] + return n > 0 + + def run_if_due(self): + """供采集循环调用:有新快讯且冷却结束才分析""" + import time as _t + if not self.llm.enabled: + return + now = _t.time() + if now - self.last_run < ANALYZE_COOLDOWN_S: + return + conn = self._conn() + try: + self._ensure_tables(conn) + if not self.has_new_news(conn): + return + self.last_run = now + self.analyze_latest(conn) + except LlmError as e: + self._emit('事件分析跳过', {'msg': str(e)[:120]}) + except Exception as e: + traceback.print_exc() + self._emit('事件分析失败', {'msg': str(e)[:120]}) + finally: + conn.close() + + # ---------- 单次分析 ---------- + + def analyze_latest(self, conn): + news = conn.execute( + "SELECT event_time, content FROM news_flash " + "WHERE event_time >= datetime('now', '-6 hours') " + "ORDER BY event_time DESC LIMIT ?", (MAX_NEWS_PER_RUN,)).fetchall() + if not news: + return None + today = datetime.now().strftime('%Y-%m-%d') + top_flow = conn.execute( + "SELECT ts_code, ROUND(main_net_in/1e8,2) y FROM money_flow WHERE trade_date=? " + "ORDER BY main_net_in DESC LIMIT 8", (today,)).fetchall() + spot = self._market_snapshot() + + evidence = [{'time': t, 'text': c[:160]} for t, c in news] + prompt_parts = ['## 最新快讯(原文)'] + prompt_parts += ['[{}] {}'.format(t, c) for t, c in news] + prompt_parts.append('\n## 今日主力资金净流入 TOP8(亿元)') + prompt_parts += ['{}: +{}亿'.format(r[0], r[1]) for r in top_flow] or ['(无)'] + prompt_parts.append('\n## 上证指数') + sh = conn.execute("SELECT trade_date, close FROM stock_daily WHERE ts_code='sh000001' " + "ORDER BY trade_date DESC LIMIT 1").fetchone() + if sh: + prompt_parts.append('{} 收盘 {}'.format(sh[0], sh[1])) + prompt_parts.append('\n## 全市场个股现价快照(节选,只能从中选股)\n') + prompt_parts.append('代码 | 名称 | 现价 | 今日涨跌% | 主力净流入(亿)') + for code, price, pct, mflow in spot[:120]: + prompt_parts.append('{} | {:.2f} | {:.2f}% | {:.2f}亿'.format(code, price, pct, mflow)) + user_prompt = '\n'.join(prompt_parts) + + card = None + for attempt in range(2): + content = self.llm.chat(SYSTEM_PROMPT, + user_prompt + ('\n\n上一次输出不是合法 JSON,请重新输出。' if attempt else '')) + card = self._parse(content, spot_map=None) + if card: + break + if not card or not card.get('has_event'): + self._mark_analyzed(conn, news) + return None + self._mark_analyzed(conn, news) + + card['stocks'] = self._verify_stocks(card.get('stocks') or []) + event = { + 'ts': datetime.now().strftime('%Y-%m-%d %H:%M:%S'), + 'kind': '事件影响分析', + 'data': { + 'event_summary': card.get('event_summary', ''), + 'direction': card.get('direction', '中性'), + 'sectors': card.get('sectors', []), + 'confidence': card.get('confidence', 0), + 'reasoning': card.get('reasoning', ''), + 'stocks': card['stocks'], + 'evidence': evidence, + }, + } + self._save(conn, event) + self.emit(event) + return event + + # ---------- grounding ---------- + + def _market_snapshot(self): + """全市场现价快照(新浪源,一次请求):[(code, price, pct, main_net_in亿)]""" + import akshare as ak + df = ak.stock_zh_a_spot() + df.columns = [str(c) for c in df.columns] + out = [] + flow = self._today_flow_map() + for _, r in df.iterrows(): + raw = str(r.get('代码', '')) + code = raw[-6:] if raw else '' + if code[:2] not in ('60', '00', '30', '68'): # 仅沪深A股 + continue + try: + price = float(r.get('最新价')) + except (TypeError, ValueError): + continue + if price <= 0: + continue + out.append((code, price, float(r.get('涨跌幅') or 0), flow.get(code, 0.0))) + return out + + def _today_flow_map(self): + today = datetime.now().strftime('%Y-%m-%d') + try: + return {code: round(v / 1e8, 2) for code, v in self._conn().execute( + "SELECT ts_code, main_net_in FROM money_flow WHERE trade_date=? AND main_net_in IS NOT NULL", + (today,))} + except sqlite3.Error: + return {} + + def _verify_stocks(self, stocks): + """校验 LLM 选出的股票必须存在于真实快照;购入区间以真实现价重算""" + spot = {code: (price, pct, flow) for code, price, pct, flow in self._market_snapshot()} + verified = [] + for s in stocks[:5]: + code = str(s.get('code', '')).zfill(6)[:6] + if code not in spot: + continue + price, pct, flow = spot[code] + try: + lo = float(s.get('buy_zone_low', price * 0.99)) + hi = float(s.get('buy_zone_high', price * 1.01)) + except (TypeError, ValueError): + lo, hi = price * 0.99, price * 1.01 + lo, hi = min(lo, hi), max(lo, hi) + try: + er = max(-15.0, min(15.0, float(s.get('expected_return_pct', 0)))) + except (TypeError, ValueError): + er = 0.0 + try: + stars = max(1, min(5, int(s.get('stars', 3)))) + except (TypeError, ValueError): + stars = 3 + verified.append({ + 'name': s.get('name', ''), 'code': code, + 'price': round(price, 2), 'pct_today': round(pct, 2), + 'main_net_in_yi': round(flow, 2), + 'buy_zone': [round(lo, 2), round(hi, 2)], + 'expected_return_pct': er, 'stars': stars, + 'reason': s.get('reason', ''), + }) + verified.sort(key=lambda x: -x['stars']) + return verified + + def _save(self, conn, event): + conn.execute( + "INSERT INTO event_impact(ts,event_summary,direction,sectors,confidence," + "reasoning,stocks,evidence) VALUES (?,?,?,?,?,?,?,?)", + (event['ts'], event['data']['event_summary'], event['data']['direction'], + json.dumps(event['data'].get('sectors', []), ensure_ascii=False), + event['data']['confidence'], event['data']['reasoning'], + json.dumps(event['data']['stocks'], ensure_ascii=False), + json.dumps(event['data'].get('evidence', []), ensure_ascii=False))) + conn.commit() + + def _mark_analyzed(self, conn, news): + now = datetime.now().strftime('%Y-%m-%d %H:%M:%S') + conn.executemany("INSERT OR REPLACE INTO news_analyzed(event_time,content,analyzed_at) " + "VALUES (?,?,?)", [(t, c, now) for t, c in news]) + conn.commit() + + def _parse(self, content, spot_map=None): + import re + m = re.search(r'\{.*\}', content, re.S) + if not m: + return None + try: + return json.loads(m.group(0)) + except json.JSONDecodeError: + return None + + def history(self, limit=20): + conn = self._conn() + rows = conn.execute("SELECT ts,event_summary,direction,sectors,confidence,stocks,evidence " + "FROM event_impact ORDER BY id DESC LIMIT ?", (limit,)).fetchall() + conn.close() + out = [] + for ts, summary, direction, sectors, conf, stocks, evidence in rows: + out.append({'ts': ts, 'kind': '事件影响分析', + 'data': {'event_summary': summary, 'direction': direction, + 'sectors': json.loads(sectors or '[]'), + 'stocks': json.loads(stocks or '[]'), + 'evidence': json.loads(evidence or '[]'), + 'confidence': conf}}) + return out diff --git a/src/analysis/llm_client.py b/src/analysis/llm_client.py new file mode 100644 index 0000000..47a8abf --- /dev/null +++ b/src/analysis/llm_client.py @@ -0,0 +1,76 @@ +# -*- coding: utf-8 -*- +""" +LLM 客户端(OpenAI 兼容 /chat/completions)。 +- base_url / model / key 由环境变量配置(默认智谱 GLM) +- 出站安全:仅 http/https、显式拒绝 localhost、解析 IP 拒绝环回/私有/保留段、 + 禁用重定向(防 DNS rebinding 绕过),值全部走 JSON 序列化,密钥不落日志 +""" +import ipaddress +import json +import os +from urllib.parse import urlparse + +import requests + + +class LlmError(Exception): + pass + + +def _validated_url(base_url: str) -> str: + u = urlparse(base_url) + if u.scheme not in ('http', 'https'): + raise LlmError('LLM base_url 仅允许 http/https') + host = u.hostname or '' + if not host or host.lower() in ('localhost', 'localhost.localdomain'): + raise LlmError('LLM base_url 拒绝 localhost') + port = u.port or (443 if u.scheme == 'https' else 80) + try: + infos = socket.getaddrinfo(host, port) + except socket.gaierror as e: + raise LlmError('LLM base_url 域名解析失败: {}'.format(host)) + for info in infos: + ip = ipaddress.ip_address(info[4][0]) + if (ip.is_loopback or ip.is_private or ip.is_link_local or ip.is_reserved + or ip.is_multicast or ip.is_unspecified): + raise LlmError('LLM base_url 拒绝非公网地址: {}'.format(ip)) + return '{}://{}{}'.format(u.scheme, u.netloc, u.path) + + +class LlmClient: + + def __init__(self): + self.api_key = os.environ.get('JQUANT_LLM_API_KEY', '').strip() + self.base_url = (os.environ.get('JQUANT_LLM_BASE_URL', '').strip() + or 'https://open.bigmodel.cn/api/paas/v4') + self.model = os.environ.get('JQUANT_LLM_MODEL', '').strip() or 'glm-4-flash' + self.temperature = float(os.environ.get('JQUANT_LLM_TEMPERATURE', '0.2')) + + @property + def enabled(self) -> bool: + return bool(self.api_key) + + def chat(self, system_prompt: str, user_prompt: str) -> str: + if not self.enabled: + raise LlmError('未配置 LLM API Key(JQUANT_LLM_API_KEY),请在服务环境变量中设置') + url = _validated_url(self.base_url.rstrip('/')) + '/chat/completions' + body = { + 'model': self.model, + 'temperature': self.temperature, + 'messages': [ + {'role': 'system', 'content': system_prompt}, + {'role': 'user', 'content': user_prompt}, + ], + } + # 校验与请求紧邻;禁重定向防 DNS rebinding 绕过 IP 校验 + r = requests.post(url, json=body, timeout=90, allow_redirects=False, + headers={'Authorization': 'Bearer ' + self.api_key, + 'Content-Type': 'application/json'}) + if r.status_code in (301, 302, 303, 307, 308): + raise LlmError('LLM 端点发生重定向,已拒绝(防 SSRF 绕过)') + if r.status_code != 200: + raise LlmError('LLM HTTP {}: {}'.format(r.status_code, r.text[:200])) + content = r.json().get('choices', [{}])[0].get('message', {}).get('content') + if not content: + raise LlmError('LLM 响应缺少 content') + return content diff --git a/src/realtime/collector.py b/src/realtime/collector.py index b263907..2332444 100644 --- a/src/realtime/collector.py +++ b/src/realtime/collector.py @@ -45,6 +45,8 @@ class TimelineCollector: FEED.parent.mkdir(parents=True, exist_ok=True) # 事件流水目录 self.priors = {} self._load_priors() + from src.analysis.event_impact import EventImpactService + self.impact = EventImpactService(emit=self.emit, db_path=self.db_path) # ---------- 基础 ---------- @@ -177,7 +179,7 @@ class TimelineCollector: conn.execute( "INSERT INTO news_flash(content, event_time, source, inserted_at) " "VALUES (?,?,?,datetime('now','localtime'))", (summary, ts, '东财7x24')) - fresh.append((ts, summary[:60])) + fresh.append((ts, summary[:200])) conn.commit() except Exception as e: self._emit('采集错误', {'msg': str(e)[:120]}) @@ -236,6 +238,7 @@ class TimelineCollector: self._emit('盘中点位', {'上证': spot['price'], '涨跌幅%': spot['pct']}) for h in self.news_alerts(conn): self._emit('快讯关键词告警', h) + self.impact.run_if_due() if hm >= 1510 and not self.state.get('postmarket_done'): n = self.save_lhb_today(conn) diff --git a/src/web/index.html b/src/web/index.html index 85e9a84..8a57bb8 100644 --- a/src/web/index.html +++ b/src/web/index.html @@ -49,6 +49,16 @@ .k.盘前隔夜预估, .k.收盘日报 { color: var(--down); } .d { flex: 1; word-break: break-all; color: #c3cad8; } .empty { text-align: center; color: var(--text2); padding: 40px 0; } + .impact { margin-top: 8px; border-top: 1px dashed var(--border); padding-top: 8px; } + .impact table { width: 100%; border-collapse: collapse; font-size: 12px; } + .impact th { color: var(--text2); text-align: left; padding: 3px 8px; font-weight: 600; } + .impact td { padding: 4px 8px; border-top: 1px solid var(--border); font-variant-numeric: tabular-nums; } + .stars { color: #ffb74d; letter-spacing: 1px; } + .dir-bull { color: var(--up); font-weight: 700; } + .dir-bear { color: var(--down); font-weight: 700; } + .dir-neutral { color: var(--text2); } + .conf { font-size: 11px; color: var(--text2); } + .disclaimer { margin-top: 6px; font-size: 11px; color: #6b7280; }
@@ -65,6 +75,7 @@ +| 标的 | 代码 | 建议购入区间 | 预计收益(5日,推测) | 推荐指数 | 简易原因 |
|---|