From f20fe2f4f0e095389d526709412add640f0701b5 Mon Sep 17 00:00:00 2001 From: hz4th_coder Date: Wed, 9 Sep 2026 15:59:35 +0800 Subject: [PATCH] =?UTF-8?q?v1.9.0:=20=E2=91=A0=E8=87=AA=E5=8A=A8=E5=8C=96?= =?UTF-8?q?=E5=BC=80=E5=85=B3=E5=8D=B3=E6=97=B6=E4=BF=9D=E5=AD=98(?= =?UTF-8?q?=E5=88=87=E6=8D=A2=E5=8D=B3=E7=94=9F=E6=95=88=E4=B8=8D=E5=9B=9E?= =?UTF-8?q?=E8=90=BD)=20=E2=91=A1=E6=8C=81=E4=BB=93=E8=B7=9F=E8=B8=AA?= =?UTF-8?q?=E6=8E=A8=E9=80=81=E6=96=B9=E5=BC=8F(=E8=81=9A=E5=90=88?= =?UTF-8?q?=E4=B8=80=E8=B5=B7=E5=8F=91=E9=BB=98=E8=AE=A4/=E5=88=86?= =?UTF-8?q?=E6=95=A3=E5=8F=91)=20=E2=91=A2=E8=A1=8C=E6=83=85=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E5=A2=9E=E9=87=8F=E5=BB=B6=E7=BB=AD=E5=88=B0=E6=9C=80?= =?UTF-8?q?=E6=96=B0=E4=BA=A4=E6=98=93=E6=97=A5(=E6=AF=8F=E5=B0=8F?= =?UTF-8?q?=E6=97=B6=E8=87=AA=E5=8A=A8+=E6=89=8B=E5=8A=A8=E6=8C=89?= =?UTF-8?q?=E9=92=AE)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 1 + app.py | 54 +++++++++++++++++ engine/agent.py | 102 ++++++++++++++++++++++++++++---- seed_data.py | 119 +++++++++++++++++++++++++++++++++++++- settings.py | 2 + static/js/admin.js | 17 ++++++ static/js/automation.js | 4 +- templates/admin.html | 4 +- templates/automation.html | 46 +++++++++------ 9 files changed, 317 insertions(+), 32 deletions(-) diff --git a/README.md b/README.md index 9530e91..7601a74 100644 --- a/README.md +++ b/README.md @@ -389,6 +389,7 @@ IS_MOCK = False | v1.6.0 | 目标跟踪:概念/主题/股票(watch_targets),概念智能体分析+受益个股,持仓管理与目标管理,报告支持双类型 | | v1.7.0 | 每日定时报告:工作日 9:00 盘前 / 15:30 盘后,简版正文+详细版附件,全球市场数据,cron 定时+手动触发 | | v1.8.0 | 🌙 静默期:舆情监控/持仓跟踪/定时报告三大自动化任务均可独立配置静默时段(多段/跨午夜),命中时段不发送邮件(日志标记 quiet) | +| v1.9.0 | ①自动化页开关/参数全部**即时保存**(切换即生效,无需再点保存按钮,刷新不回弹);②持仓跟踪新增**重大变化推送方式**:聚合一起发(默认,一轮全部分析完统一发一封汇总邮件)/ 分散发(每只出结果立即单独发);③**行情数据自动延续**:新增 extend_daily 增量扩展(行情/新闻/指数从库内最新日期延续到最新交易日,幂等),服务每小时自动检查+数据管理页「📅 更新行情到最新交易日」手动触发,数据不再停留在建库日 | --- diff --git a/app.py b/app.py index 1ef8564..f01fe3d 100644 --- a/app.py +++ b/app.py @@ -619,6 +619,8 @@ def api_settings_save(): set_setting("tracking_enabled", _bool_str(track["tracking_enabled"])) if "tracking_notify" in track: set_setting("tracking_notify", _bool_str(track["tracking_notify"])) + if "tracking_push_mode" in track: + set_setting("tracking_push_mode", str(track["tracking_push_mode"]).strip().lower()) # 静默期(各自动化任务独立配置) quiet = body.get("quiet") or {} for prefix in ("monitor", "tracking", "report"): @@ -786,6 +788,8 @@ def api_admin_stats(): tables = ("stocks", "stock_daily", "news", "institutions", "inst_ratings", "fund_holdings", "watchlist", "analysis_cache", "analysis_history", "strategy_backtests", "market_index") + latest_daily = query_one("SELECT MAX(date) d FROM stock_daily") + latest_news = query_one("SELECT MAX(publish_date) d FROM news") return jsonify({ "tables": {t: table_count(t) for t in tables}, "vector": { @@ -794,6 +798,8 @@ def api_admin_stats(): }, "is_mock": IS_MOCK, "db": "stock_advisor.db", + "latest_daily": (latest_daily or {}).get("d"), + "latest_news": (latest_news or {}).get("d"), }) @@ -814,6 +820,29 @@ def api_admin_reseed(): return jsonify({"ok": True, "msg": "重灌任务已启动,可在数据管理页刷新查看进度"}) +@app.route("/api/admin/update-data", methods=["POST"]) +def api_admin_update_data(): + """把行情/新闻/指数增量扩展到最新交易日(模拟数据延续),后台执行""" + import threading + + def run(): + try: + from seed_data import extend_daily + with open(os.path.join(LOG_DIR, "data_update.log"), "a") as f: + f.write(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] 开始更新数据\n") + res = extend_daily() + with open(os.path.join(LOG_DIR, "data_update.log"), "a") as f: + f.write(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] 完成: {res}\n") + log.info("data extend done: %s", res) + except Exception as e: + log.exception("data extend fail") + with open(os.path.join(LOG_DIR, "data_update.log"), "a") as f: + f.write(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] 失败: {e}\n") + + threading.Thread(target=run, daemon=True).start() + return jsonify({"ok": True, "msg": "行情更新已启动(增量扩展到最新交易日),稍后刷新数据管理页查看"}) + + @app.route("/api/admin/healthcheck") def api_admin_healthcheck(): """外部依赖连通性检查""" @@ -835,10 +864,35 @@ def api_admin_healthcheck(): if __name__ == "__main__": + import threading init_db() from engine.notifier import start_monitor from engine.agent import start_tracking start_monitor() start_tracking() + # 行情自动更新器:每小时检查,若行情最新日期落后于最新交易日则增量扩展(模拟数据延续) + class DataUpdater(threading.Thread): + def __init__(self): + super().__init__(daemon=True, name="data-updater") + self._stop = threading.Event() + + def run(self): + log.info("行情自动更新器启动(每小时检查)") + while not self._stop.is_set(): + try: + import datetime as _dt + last = query_one("SELECT MAX(date) d FROM stock_daily") + today = _dt.date.today() + while today.weekday() >= 5: + today -= _dt.timedelta(days=1) + if not last or last["d"] < today.isoformat(): + from seed_data import extend_daily + res = extend_daily() + log.info("行情自动扩展: %s", res) + except Exception as e: + log.warning("data updater error: %s", e) + self._stop.wait(3600) + + DataUpdater().start() print(f"✅ {SERVICE_NAME} 启动: http://0.0.0.0:{SERVICE_PORT}") app.run(host=SERVICE_HOST, port=SERVICE_PORT, threaded=True) diff --git a/engine/agent.py b/engine/agent.py index fddc6f3..52681c1 100644 --- a/engine/agent.py +++ b/engine/agent.py @@ -222,8 +222,8 @@ def strip_json_block(text): # ===================================================================== 执行 -def track_stock(code, focus=""): - """执行一次跟踪,返回 {ok, report_id, meta, ...}""" +def track_stock(code, focus="", notify=True): + """执行一次跟踪,返回 {ok, report_id, meta, ...}。notify=False 时不单独发邮件(聚合推送模式)""" with _track_lock: seg = collect_chain(code) if not seg: @@ -267,7 +267,8 @@ def track_stock(code, focus=""): (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"] - _notify_if_significant(rid, s, meta) + if notify: + _notify_if_significant(rid, s, meta) return {"ok": True, "report_id": rid, "meta": meta} except Exception as e: log.exception("track %s fail", code) @@ -323,8 +324,9 @@ def _notify_if_significant(rid, stock, meta): # ===================================================================== 概念/主题跟踪 -def track_concept(name, keywords="", focus=""): - """跟踪一个概念/主题:采集相关新闻 → 受益个股梳理 → LLM 深度分析 → 影响度判定""" +def track_concept(name, keywords="", focus="", notify=True): + """跟踪一个概念/主题:采集相关新闻 → 受益个股梳理 → LLM 深度分析 → 影响度判定。 + notify=False 时不单独发邮件(聚合推送模式)""" kw = keywords.strip() or name kw_list = [k.strip() for k in kw.replace(",", ",").split(",") if k.strip()] kws = " ".join(kw_list) @@ -411,7 +413,8 @@ def track_concept(name, keywords="", focus=""): (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) + if notify: + _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) @@ -437,9 +440,69 @@ def delete_target(tid): return {"ok": True} +def _notify_aggregate(items): + """聚合推送:一轮跟踪全部完成后,统一发一封汇总邮件(默认推送方式)""" + if not items: + return + cfg = tracking_config() + try: + if not cfg["notify"]: + return + qc = quiet_config("tracking") + if qc["enabled"] and in_quiet_period(qc["ranges"]): + log.info("tracking aggregate notify suppressed by quiet period (%d items)", len(items)) + return + from engine.notifier import send_email + mc = mail_config() + rows = "" + for it in items: + meta = it.get("meta") or {} + is_concept = it.get("type") == "concept" + impact = int(meta.get("impact_score", 0)) + level = "🔴 重大" if impact >= 65 else "🟡 关注" + name = escape_html(it.get("name", "")) + code = it.get("code", "") + rows += ( + f"" + f"{'🎯 概念' if is_concept else '💼 股票'}" + f"{name}" + f"{'
' + str(code) + '
' if not is_concept and code else ''}" + f"{meta.get('change_kind', '')}" + f"{impact}/100 ({level})" + f"{escape_html(meta.get('summary', ''))}" + f"" + ) + send_email( + f"[持仓跟踪] 聚合报告:{len(items)} 个目标出现重大动态(影响度≥{int(cfg['impact_threshold'])})", + f""" +
+
🧭 持仓跟踪聚合报告 · {time.strftime('%Y-%m-%d %H:%M')}
+
+

本轮共跟踪 {items[0].get('_total', len(items))} 个目标,其中 {len(items)} 个达到重大变化判定阈值(影响度 ≥ {int(cfg['impact_threshold'])}):

+ + +{rows} +
类型目标性质影响度摘要
+

详细分析报告请在系统「🧭 持仓跟踪智能体 → 最近跟踪报告」中查看。

+
""", + cfg=mc) + set_tracking_state(last_alert=int(items[0]["meta"].get("impact_score", 0))) + log.info("tracking aggregate notify sent: %d items", len(items)) + except Exception as e: + log.warning("tracking aggregate notify fail: %s", e) + + +def escape_html(s): + return str(s or "").replace("&", "&").replace("<", "<").replace(">", ">").replace('"', """) + + # ===================================================================== 批量与调度 def track_all(progress=None): - """跟踪全部目标:持仓(自选股)+ 目标(概念/主题/股票)""" + """跟踪全部目标:持仓(自选股)+ 目标(概念/主题/股票)。 + 推送方式(tracking_push_mode): + aggregate(默认) —— 一轮全部分析完后,汇总重大变化统一发一封邮件; + scattered —— 每只股票分析完立即单独发邮件。 + """ global _cycle_running if _cycle_running: return {"tracked": 0, "msg": "已有跟踪任务进行中,请稍后再试"} @@ -465,16 +528,33 @@ def track_all(progress=None): uniq.append((typ, key, name)) if not uniq: return {"tracked": 0, "msg": "暂无跟踪目标:请添加持仓/自选股或概念主题目标"} + cfg = tracking_config() + scattered = cfg.get("push_mode", "aggregate") == "scattered" + significant = [] # 聚合推送模式收集重大变化 results = [] for i, (typ, key, name) in enumerate(uniq): - if typ == "stock": - r = track_stock(key) - else: - r = track_concept(key, name) + try: + if typ == "stock": + r = track_stock(key, notify=scattered) + else: + r = track_concept(key, name, notify=scattered) + except Exception as e: + r = {"error": str(e)} results.append({"type": typ, "name": name, **r}) + if r.get("ok"): + meta = r.get("meta") or {} + if int(meta.get("impact_score", 0)) >= int(cfg["impact_threshold"]): + significant.append({"type": typ, "name": name, "code": key, + "meta": meta, "_total": len(uniq)}) set_tracking_state(last_run=time.strftime("%Y-%m-%d %H:%M:%S"), last_stock=name) if progress: progress(i + 1, len(uniq)) + # 聚合推送:一轮分析完统一发一封汇总邮件 + if not scattered and significant: + try: + _notify_aggregate(significant) + except Exception as e: + log.warning("aggregate notify error: %s", e) return {"tracked": len(results), "results": results} finally: _cycle_running = False diff --git a/seed_data.py b/seed_data.py index c033e98..6cf6ee4 100644 --- a/seed_data.py +++ b/seed_data.py @@ -21,7 +21,7 @@ import random import sys from config import (CHROMA_NEWS_COLLECTION, CHROMA_PROFILE_COLLECTION, LOG_DIR) -from database import init_db, executemany, query_one, wipe_all +from database import init_db, executemany, query_one, query, execute, wipe_all from rag import vector_store as vs random.seed(42) @@ -427,6 +427,123 @@ def build_vectors(news, stocks): print(f" 概况索引条数: {vs.collection_count(CHROMA_PROFILE_COLLECTION)}") +# ===================================================================== 增量更新 +def extend_daily(target_date=None): + """把模拟行情/指数/新闻从库里最新日期增量扩展到最新交易日(幂等,可重复调用)。 + 返回 (起始日期, 结束日期, 新增交易日数, 新增新闻数) 或 (None, None, 0, 0)。 + """ + init_db() + last = query_one("SELECT MAX(date) d FROM stock_daily") + last_date = last["d"] if last else None + today = dt.date.today() + if target_date is None: + target_date = today + if isinstance(target_date, str): + target_date = dt.date.fromisoformat(target_date) + # 最新交易日:今天若是周末则回退到周五 + while target_date.weekday() >= 5: + target_date -= dt.timedelta(days=1) + if last_date: + start = dt.date.fromisoformat(last_date) + dt.timedelta(days=1) + else: + start = target_date - dt.timedelta(days=179) + dates = [] + d = start + while d <= target_date: + if d.weekday() < 5: + dates.append(d.isoformat()) + d += dt.timedelta(days=1) + if not dates: + return None, None, 0, 0 + print(f">>> 增量扩展行情 {dates[0]} ~ {dates[-1]}({len(dates)} 个交易日)...") + + stocks = query("SELECT code,name,industry,board,total_shares,float_shares FROM stocks") + daily = [] + index = {} + price = {} + for s in stocks: + rows = query("SELECT date,close,volume FROM stock_daily WHERE code=? ORDER BY date DESC LIMIT 30", + (s["code"],)) + if not rows: + continue + last_close = rows[0]["close"] + base_vol = rows[0]["volume"] or 1 + closes = [r["close"] for r in reversed(rows)] + rets = [(closes[i + 1] - closes[i]) / closes[i] for i in range(len(closes) - 1)] + mean_r = sum(rets) / max(len(rets), 1) + vol = (sum((r - mean_r) ** 2 for r in rets) / max(len(rets), 1)) ** 0.5 if rets else 0.02 + vol = max(0.005, min(0.05, vol)) + drift = 0.0002 + p = prev = last_close + for dd in dates: + r = random.gauss(drift, vol) + if random.random() < 0.02: + r += random.gauss(0, vol * 1.6) + prev = p + p = max(0.5, p * (1 + r)) + open_p = prev * (1 + random.gauss(0, vol * 0.5)) + high = max(open_p, p) * (1 + abs(random.gauss(0, vol * 0.35))) + low = min(open_p, p) * (1 - abs(random.gauss(0, vol * 0.35))) + volume = base_vol * (1 + 1.5 * abs(r) / max(vol, 1e-6)) * random.uniform(0.6, 1.4) + amount = volume * (open_p + p) / 2 + chg = (p - prev) / prev * 100 + daily.append((s["code"], dd, round(open_p, 2), round(high, 2), round(low, 2), + round(p, 2), round(volume, 0), round(amount, 0), round(chg, 2))) + price[s["code"]] = p + + # 指数延续 + idx_rows = query("SELECT date,sh,sz,cy FROM market_index ORDER BY date DESC LIMIT 1") + if idx_rows: + sh, sz, cy = idx_rows[0]["sh"], idx_rows[0]["sz"], idx_rows[0]["cy"] + else: + sh, sz, cy = 3245.0, 10580.0, 2120.0 + for dd in dates: + sh_r = sum(random.gauss(0.0004, 0.008) for _ in range(6)) / 6 + sz_r = sh_r + random.gauss(0, 0.004) + cy_r = sh_r + random.gauss(0, 0.006) + sh *= (1 + sh_r); sz *= (1 + sz_r); cy *= (1 + cy_r) + index[dd] = {"sh": round(sh, 2), "sz": round(sz, 2), "cy": round(cy, 2)} + + executemany( + "INSERT OR REPLACE INTO stock_daily(code,date,open,high,low,close,volume,amount,change_pct) " + "VALUES(?,?,?,?,?,?,?,?,?)", daily) + executemany("INSERT OR REPLACE INTO market_index(date,sh,sz,cy) VALUES(?,?,?,?)", + [(d, v["sh"], v["sz"], v["cy"]) for d, v in index.items()]) + gm = _gen_global(dates) + executemany("INSERT OR REPLACE INTO global_markets(date,data) VALUES(?,?)", + [(d, json.dumps(v, ensure_ascii=False)) for d, v in gm.items()]) + # 回填市值 + for code in price: + execute("UPDATE stocks SET market_cap=ROUND((SELECT close FROM stock_daily WHERE code=? ORDER BY date DESC LIMIT 1)*total_shares,2) WHERE code=?", + (code, code)) + + # 新增新闻(仅保留落在新增日期内的) + news = gen_news(dates, price) + date_set = set(dates) + new_news = [n for n in news if n["publish_date"] in date_set] + if new_news: + executemany( + "INSERT INTO news(title,content,source,category,publish_date,related_stocks,sentiment,is_positive) " + "VALUES(?,?,?,?,?,?,?,?)", + [(n["title"], n["content"], n["source"], n["category"], n["publish_date"], + n["related"], n["sentiment"], n["is_positive"]) for n in new_news]) + # 向量追加 + ids, docs, metas = [], [], [] + for n in new_news: + for code in n["related"].split(","): + ids.append(f"news-{n['publish_date']}-{n['title']}-{code}") + docs.append(f"{n['title']}\n{n['content']}") + metas.append({"code": code, "title": n["title"], "date": n["publish_date"], + "category": n["category"], "sentiment": n["sentiment"], "news_id": 0}) + try: + for i in range(0, len(ids), 16): + vs.add_documents(ids[i:i + 16], docs[i:i + 16], metas[i:i + 16], CHROMA_NEWS_COLLECTION) + except Exception as e: + print("新闻向量追加失败(可忽略,下次重灌会重建):", e) + print(f">>> 新增新闻 {len(new_news)} 条(向量已追加)") + return dates[0], dates[-1], len(dates), len(new_news) + + def main(): parser = argparse.ArgumentParser() parser.add_argument("--skip-vector", action="store_true", help="跳过向量索引重建") diff --git a/settings.py b/settings.py index edc93a9..4017196 100644 --- a/settings.py +++ b/settings.py @@ -115,6 +115,8 @@ def tracking_config(): "interval_min": max(15, int(get_setting("tracking_interval", TRACKING_DEFAULTS["tracking_interval"]))), "notify": _truthy(get_setting("tracking_notify", TRACKING_DEFAULTS["tracking_notify"])), "impact_threshold": float(get_setting("tracking_impact_threshold", TRACKING_DEFAULTS["tracking_impact_threshold"])), + # 推送方式: aggregate=聚合一起发(默认, 一轮全部分析完统一发一封汇总) / scattered=分散发(每只出结果立即单独发) + "push_mode": (get_setting("tracking_push_mode", "aggregate") or "aggregate").strip().lower(), } diff --git a/static/js/admin.js b/static/js/admin.js index 81c3fe1..4345ef0 100644 --- a/static/js/admin.js +++ b/static/js/admin.js @@ -22,6 +22,12 @@ async function refreshStats() { 🏢 公司概况索引 (stock_profiles_v1) ${d.vector.profiles} +
+ 行情最新日期${d.latest_daily || '—'} +
+
+ 新闻最新日期${d.latest_news || '—'} +
数据模式 ${d.is_mock ? '模拟数据' : '真实数据'} @@ -98,6 +104,17 @@ async function reseed() { } } +async function updateData() { + try { + const r = await api('/api/admin/update-data', { method: 'POST' }); + toast(r.msg || '行情更新已启动'); + setTimeout(refreshStats, 8000); + setTimeout(refreshStats, 20000); + } catch (e) { + toast('启动失败:' + e.message); + } +} + async function health() { $('#healthBox').innerHTML = '
检查中…
'; try { diff --git a/static/js/automation.js b/static/js/automation.js index 96c2e39..f2a797a 100644 --- a/static/js/automation.js +++ b/static/js/automation.js @@ -30,6 +30,7 @@ async function loadAuto() { $('#trackingInterval').value = tr.interval_min; $('#trackingImpact').value = tr.impact_threshold; $('#trackingNotify').checked = !!tr.notify; + $('#trackingPushMode').value = (tr.push_mode === 'scattered') ? 'scattered' : 'aggregate'; renderTrackingState(tr); // 静默期(跟踪) const tq = (autoSettings.quiet && autoSettings.quiet.tracking) || {}; @@ -100,7 +101,8 @@ async function saveTracking() { tracking_enabled: $('#trackingEnabled').checked, tracking_interval: $('#trackingInterval').value, tracking_impact_threshold: $('#trackingImpact').value, - tracking_notify: $('#trackingNotify').checked + tracking_notify: $('#trackingNotify').checked, + tracking_push_mode: $('#trackingPushMode').value }, quiet: { tracking: collectTrackingQuiet() } }) }); toast(r.msg || '已保存'); } catch (e) { toast('保存失败:' + e.message); } diff --git a/templates/admin.html b/templates/admin.html index 9ded37a..bba491a 100644 --- a/templates/admin.html +++ b/templates/admin.html @@ -47,13 +47,15 @@
系统维护
+
- 当前为模拟数据,用于功能演示与系统验证。接入真实数据后(行情/新闻/机构),修改 config.py 中的 IS_MOCK=False,并替换 seed_data.py 为真实数据源即可,其余分析/检索逻辑无需改动。
+ 当前为模拟数据,用于功能演示与系统验证。📅 更新行情会把模拟行情/新闻/指数从库内最新日期增量延续到最新交易日(服务每小时自动检查一次,无需手动)。
+ 接入真实数据后(行情/新闻/机构),修改 config.py 中的 IS_MOCK=False,并替换 seed_data.py 为真实数据源即可,其余分析/检索逻辑无需改动。
技术栈:Python Flask + SQLite + Chroma 向量库 + bge-large-zh 语义检索 + DeepSeek 大模型(RAG 增强研报)。
diff --git a/templates/automation.html b/templates/automation.html index fd9432a..42adac5 100644 --- a/templates/automation.html +++ b/templates/automation.html @@ -12,30 +12,30 @@
定期扫描最新新闻,重要信息邮件通知
- +
-
-
-
+
+
+
- - - - - + + + + +
-
+
🌙 静默期 —— 某些时段不发送邮件通知
- +
-
+
开启后,落在静默时段内的扫描将不发送邮件(通知日志标记为“静默”,水位正常推进,不会在结束后补发积压)。
@@ -59,22 +59,32 @@