From 40a77cad2c211e844fbc666c004bfe19ff77645a Mon Sep 17 00:00:00 2001 From: Yuzhiran Dev Date: Wed, 20 May 2026 19:07:35 +0800 Subject: [PATCH] =?UTF-8?q?Phase2:=20=E7=9C=9F=E5=AE=9E=E7=83=AD=E7=82=B9?= =?UTF-8?q?=E6=95=B0=E6=8D=AE=E6=8E=A5=E5=85=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - trends.py新增百度/微博/知乎实时热搜API抓取,LLM为fallback - 新增 source 字段标记数据来源 - scheduler.py新增 scheduled_fetch_trends 每日03:00定时刷新 - system.py新增 POST /api/system/trends/run 手动触发端点 - system.py modules/status 加入热点趋势模块 - index.html triggerModule 加入 trends 触发按钮 --- platform/backend/app/api/system.py | 18 ++ platform/backend/app/core/scheduler.py | 31 +++- platform/frontend/index.html | 1 + scripts/trends.py | 231 ++++++++++++++++++++++--- 4 files changed, 258 insertions(+), 23 deletions(-) diff --git a/platform/backend/app/api/system.py b/platform/backend/app/api/system.py index e4ad49c..2962b9e 100644 --- a/platform/backend/app/api/system.py +++ b/platform/backend/app/api/system.py @@ -176,6 +176,23 @@ def trigger_metrics_sync(): except Exception as e: raise HTTPException(status_code=500, detail=str(e)) +@router.post("/trends/run") +def trigger_trends_refresh(): + """手动刷新热点趋势数据""" + try: + import subprocess, sys as sys_mod + from pathlib import Path + scripts_dir = Path(__file__).parent.parent.parent.parent / "scripts" + result = subprocess.run( + [sys_mod.executable, str(scripts_dir / "trends.py")], + capture_output=True, text=True, timeout=120 + ) + if result.returncode != 0: + raise Exception(result.stderr[-500:]) + return {"message": "热点趋势已刷新", "output": result.stdout.strip()} + except Exception as e: + raise HTTPException(status_code=500, detail=str(e)) + @router.get("/automation/topics") def list_automation_topics(db: Session = Depends(get_db), current_user=Depends(get_current_user)): try: @@ -221,6 +238,7 @@ def get_modules_status(): today_str = date.today().isoformat() log_based: dict = { "scheduled_collect": {"name": "📡 内容采集", "log": LOGS_DIR / f"collector_{today_str}.log"}, + "scheduled_fetch_trends": {"name": "🔥 热点趋势", "log": LOGS_DIR / f"trends_{today_str}.log"}, "scheduled_generate": {"name": "🤖 内容创作", "log": LOGS_DIR / f"creator_{today_str}.log"}, "scheduled_optimize": {"name": "🔍 合规审查", "log": LOGS_DIR / f"optimizer_{today_str}.log"}, "scheduled_optimize_sources": {"name": "📡 信息源优化", "log": LOGS_DIR / f"optimizer_sources_{today_str}.log"}, diff --git a/platform/backend/app/core/scheduler.py b/platform/backend/app/core/scheduler.py index 1a5f917..448133a 100644 --- a/platform/backend/app/core/scheduler.py +++ b/platform/backend/app/core/scheduler.py @@ -3,8 +3,10 @@ 基于 APScheduler,支持在 FastAPI 生命周期内运行定时任务 """ import os +import sys import logging from datetime import datetime +from pathlib import Path from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger from .generator import run_creator @@ -28,6 +30,14 @@ class TaskScheduler: CronTrigger(hour=1, minute=30), id='scheduled_collect', ) + self.scheduler.add_job( + self._run_fetch_trends, + CronTrigger(hour=3, minute=0), + id='scheduled_fetch_trends', + replace_existing=True, + max_instances=1, + coalesce=True + ) self.scheduler.add_job( self._run_generate, CronTrigger(hour=3, minute=30), @@ -62,12 +72,31 @@ class TaskScheduler: ) self.scheduler.start() self._started = True - logger.info("Scheduler started: 01:30 collect, 03:30 generate, 04:30 review, 05:00 optimize_sources, 06:00 metrics_sync") + logger.info("Scheduler started: 01:30 collect, 03:00 trends, 03:30 generate, 04:30 review, 05:00 optimize_sources, 06:00 metrics_sync") def shutdown(self): if self.scheduler.running: self.scheduler.shutdown() logger.info("Scheduler shut down") + def _run_fetch_trends(self): + """定时刷新热点趋势(百度/微博/知乎实时热搜 + LLM补充)""" + try: + logger.info("[Scheduled] Fetching hot trends...") + import subprocess + result = subprocess.run( + [sys.executable, str(Path(__file__).parent.parent.parent.parent / "scripts" / "trends.py")], + capture_output=True, text=True, timeout=120 + ) + if result.returncode == 0: + for line in result.stdout.strip().split("\n"): + if line.strip(): + logger.info("[Trends] %s", line.strip()) + logger.info("[Scheduled] Trends refreshed successfully") + else: + logger.warning("[Scheduled] Trends refresh failed: %s", result.stderr[-500:]) + except Exception as e: + logger.exception("[Scheduled] Trends refresh error: %s", e) + def _run_generate(self): try: logger.info("[Scheduled] Starting content generation...") diff --git a/platform/frontend/index.html b/platform/frontend/index.html index f73464a..17a3e60 100644 --- a/platform/frontend/index.html +++ b/platform/frontend/index.html @@ -285,6 +285,7 @@ this.runningModule = modId; const endpoints = { scheduled_collect: '/api/system/collect/run', + scheduled_fetch_trends: '/api/system/trends/run', scheduled_generate: '/api/system/generate/run', scheduled_optimize: '/api/system/review/run', scheduled_optimize_sources: '/api/system/optimize-sources/run', diff --git a/scripts/trends.py b/scripts/trends.py index df300e5..5d2811d 100644 --- a/scripts/trends.py +++ b/scripts/trends.py @@ -1,13 +1,13 @@ #!/usr/bin/env python3 """ 热点趋势感知模块 -- LLM 生成当前领域热点话题 -- 可扩展接入外部热搜 API +- 实时热搜API(百度/微博/知乎) +- LLM 生成作为 fallback """ - -import json, datetime, logging, sys +import json, datetime, logging, sys, re from pathlib import Path from typing import List, Dict +import requests PROJECT_ROOT = Path(__file__).parent.parent sys.path.insert(0, str(PROJECT_ROOT)) @@ -28,25 +28,180 @@ logger = logging.getLogger(__name__) DOMAINS = ["远程工作", "AI工具", "可持续生活", "知识管理", "数字生活", "科技人文"] +UA = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36" + +_KEYWORD_DOMAIN_MAP = [ + (r"AI|人工智能|大模型|GPT|机器学习|深度学习|聊天机器人|LLM", "AI工具"), + (r"远程|居家办公|自由职业|数字游民|远程协作", "远程工作"), + (r"可持续|环保|低碳|绿色|碳中和|循环|零浪费|垃圾分类|节能", "可持续生活"), + (r"知识管理|笔记|Obsidian|Notion|第二大脑|读书|阅读", "知识管理"), + (r"数字生活|数码|手机|电脑|智能|APP|应用|软件", "数字生活"), + (r"科技|人文|教育|心理|哲学|社会学", "科技人文"), +] + + +def _guess_domain(topic: str, reason: str = "") -> str: + text = (topic + " " + reason).lower() + for pattern, domain in _KEYWORD_DOMAIN_MAP: + if re.search(pattern, text, re.IGNORECASE): + return domain + return "科技人文" + + +def fetch_weibo_hot() -> List[Dict]: + """微博热搜""" + try: + resp = requests.get("https://weibo.com/ajax/side/hotSearch", headers={"User-Agent": UA}, timeout=10) + if resp.status_code != 200: + logger.warning(f"微博热搜返回 {resp.status_code}") + return [] + data = resp.json() + realtime = (data.get("data", {}) or {}).get("realtime", []) + results = [] + for item in realtime[:30]: + word = item.get("word", "").strip() + if not word: + continue + num = item.get("raw_hot", 0) or item.get("num", 0) + results.append({ + "domain": _guess_domain(word), + "topic": word, + "reason": f"微博热搜 热度{num}", + "hot_keywords": [word], + "platform": "微博", + "seo_angle": word, + "engagement": "高" if num > 500000 else "中", + "source": "weibo", + }) + logger.info(f"微博热搜获取 {len(results)} 条") + return results + except Exception as e: + logger.warning(f"微博热搜失败: {e}") + return [] + + +def fetch_zhihu_hot() -> List[Dict]: + """知乎热榜""" + try: + resp = requests.get( + "https://www.zhihu.com/api/v3/feed/topstory/hot-lists/total", + headers={"User-Agent": UA}, + timeout=10, + ) + if resp.status_code != 200: + logger.warning(f"知乎热榜返回 {resp.status_code}") + return [] + data = resp.json() + items = data.get("data", []) + results = [] + for item in items[:20]: + target = item.get("target", {}) + title = target.get("title", "").strip() + if not title: + continue + metrics = target.get("metrics_area", {}) + if isinstance(metrics, str): + metrics = {} + detail_text = target.get("excerpt", "") or target.get("answer", "") + results.append({ + "domain": _guess_domain(title, detail_text), + "topic": title, + "reason": f"知乎热榜", + "hot_keywords": [title] + re.findall(r'[^\s,。!?、,\.]{2,6}', title)[:3], + "platform": "知乎", + "seo_angle": title, + "engagement": "高", + "source": "zhihu", + }) + logger.info(f"知乎热榜获取 {len(results)} 条") + return results + except Exception as e: + logger.warning(f"知乎热榜失败: {e}") + return [] + + +def fetch_baidu_hot() -> List[Dict]: + """百度热搜""" + try: + resp = requests.get( + "https://top.baidu.com/api/board?tab=realtime", + headers={ + "User-Agent": UA, + "Referer": "https://top.baidu.com/", + }, + timeout=10, + ) + if resp.status_code == 200: + data = resp.json() + cards = data.get("data", {}).get("cards", []) + results = [] + for card in cards: + items = card.get("content", []) + for item in items[:30]: + word = item.get("query", "").strip() or item.get("word", "").strip() + if not word: + continue + hot_score = item.get("hotScore", 0) or item.get("heat", 0) + desc = item.get("desc", "") + results.append({ + "domain": _guess_domain(word, desc), + "topic": word, + "reason": f"百度热搜 热度{hot_score}", + "hot_keywords": [word], + "platform": "百度", + "seo_angle": word, + "engagement": "高" if hot_score > 500000 else "中", + "source": "baidu", + }) + if results: + logger.info(f"百度热搜获取 {len(results)} 条") + return results + + resp2 = requests.get("https://top.baidu.com/board?tab=realtime", headers={"User-Agent": UA}, timeout=10) + if resp2.status_code == 200: + html = resp2.text + words = re.findall(r'"word":"([^"]+)"', html) + hot_scores = re.findall(r'"hotScore":(\d+)', html) + results = [] + for i, word in enumerate(words[:30]): + score = int(hot_scores[i]) if i < len(hot_scores) else 0 + results.append({ + "domain": _guess_domain(word), + "topic": word, + "reason": f"百度热搜 热度{score}", + "hot_keywords": [word], + "platform": "百度", + "seo_angle": word, + "engagement": "高" if score > 500000 else "中", + "source": "baidu", + }) + if results: + logger.info(f"百度热搜(HTML解析)获取 {len(results)} 条") + return results + logger.warning("百度热搜获取失败") + return [] + except Exception as e: + logger.warning(f"百度热搜失败: {e}") + return [] + + def fetch_llm_trends() -> List[Dict]: - prompt = f"""你是中文互联网趋势分析师,擅长发现真正有价值、能出爆款的热点话题。请列出今天(2026年5月)中文互联网上最值得创作的10个话题。 + prompt = f"""你是中文互联网趋势分析师。请列出今天(2026年5月)中文互联网上最值得创作的10个话题。 要求: 1. 覆盖领域:{', '.join(DOMAINS)} -2. 每个话题必须从「真实用户」角度出发——不是学术热门,而是普通人正在搜、在讨论的 -3. 判断依据:知乎有高赞讨论?小红书有爆款笔记趋势?微信有刷屏文章? -4. SEO价值:这个话题是否有持续搜索量,还是纯短期流量? -5. 每个话题需包含: +2. 从真实用户角度出发 +3. 每个话题需包含: - "domain": 领域 - - "topic": 话题名称(老百姓能听懂的说法) - - "reason": 为什么现在讨论这个(1句话,说人话,有具体事件/数据支撑) - - "hot_keywords": 用户真实搜索时会用的词(3-5个,包含1-2个长尾词) - - "platform": 最适合分发此话题的平台(知乎/小红书/微信/多平台) + - "topic": 话题名称 + - "reason": 为什么现在讨论这个(1句话,有具体事件/数据支撑) + - "hot_keywords": 3-5个搜索词(含1-2个长尾词) + - "platform": 最适合分发的平台(知乎/小红书/微信/多平台) - "seo_angle": 从什么角度切入能获得搜索流量(1句话) - - "engagement": 预估互动潜力(高/中/低) + - "engagement": 高/中/低 输出 JSON 数组: -[{{"domain": "领域", "topic": "话题名", "reason": "热度原因", "hot_keywords": ["词1","词2","词3"], "platform": "知乎/小红书/微信/多平台", "seo_angle": "SEO切入点", "engagement": "高/中/低"}}] +[{{"domain": "...", "topic": "...", "reason": "...", "hot_keywords": ["..."], "platform": "...", "seo_angle": "...", "engagement": "..."}}] 只输出 JSON,不要其他文字。""" try: @@ -56,11 +211,14 @@ def fetch_llm_trends() -> List[Dict]: resp = resp.split("\n", 1)[1].rsplit("\n", 1)[0] trends = json.loads(resp) if isinstance(trends, list): + for t in trends: + t["source"] = "llm" return trends except Exception as e: logger.warning(f"LLM 趋势获取失败: {e}") return [] + def save_trends(trends: List[Dict]): data = { "date": TODAY, @@ -71,6 +229,7 @@ def save_trends(trends: List[Dict]): TRENDS_FILE.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding='utf-8') logger.info(f"趋势数据已保存: {len(trends)} 条") + def load_trends() -> List[Dict]: if TRENDS_FILE.exists(): try: @@ -81,12 +240,14 @@ def load_trends() -> List[Dict]: pass return [] + def get_trending_topics(domain: str = None, top_k: int = 5) -> List[Dict]: trends = load_trends() if domain: trends = [t for t in trends if domain in t.get("domain", "") or t.get("domain", "") in domain] return trends[:top_k] + def get_trend_context(domain: str = None) -> str: trends = get_trending_topics(domain, top_k=3) if not trends: @@ -94,21 +255,47 @@ def get_trend_context(domain: str = None) -> str: lines = ["## 当前热点趋势", ""] for t in trends: keywords = ", ".join(t.get("hot_keywords", [])) - lines.append(f"- **{t['topic']}**({t.get('platform','')}):{t.get('reason','')}") + source_tag = {"weibo": "🔥", "zhihu": "📖", "baidu": "🔍", "llm": "🤖"}.get(t.get("source", ""), "") + lines.append(f"- {source_tag} **{t['topic']}**({t.get('platform','')}):{t.get('reason','')}") if keywords: lines.append(f" 热搜词:{keywords}") return "\n".join(lines) + def main(): logger.info("开始获取热点趋势...") - trends = fetch_llm_trends() - if trends: - save_trends(trends) - for t in trends: - print(f" [{t.get('domain','?')}] {t['topic']} — {t.get('platform','')}") + all_trends = [] + seen = set() + + for fetcher_name, fetcher in [("百度热搜", fetch_baidu_hot), ("微博热搜", fetch_weibo_hot), ("知乎热榜", fetch_zhihu_hot)]: + try: + items = fetcher() + for item in items: + topic = item.get("topic", "").strip() + if topic and topic not in seen: + seen.add(topic) + all_trends.append(item) + logger.info(f"{fetcher_name}: {len(items)} 条") + except Exception as e: + logger.warning(f"{fetcher_name} 失败: {e}") + + logger.info(f"实时热搜合计 {len(all_trends)} 条,尝试LLM补充...") + llm_trends = fetch_llm_trends() + for item in llm_trends: + topic = item.get("topic", "").strip() + if topic and topic not in seen: + seen.add(topic) + all_trends.append(item) + + if all_trends: + save_trends(all_trends) + for t in all_trends[:15]: + s = {"weibo": "微", "zhihu": "知", "baidu": "百", "llm": "AI"}.get(t.get("source", ""), "?") + print(f" [{s}][{t.get('domain','?')}] {t['topic']}") + print(f"完成,共 {len(all_trends)} 条(实时 {sum(1 for t in all_trends if t.get('source') != 'llm')} 条,LLM {sum(1 for t in all_trends if t.get('source') == 'llm')} 条)") else: print("未获取到趋势数据") - print(f"完成,共 {len(trends)} 条") + if __name__ == "__main__": main()