diff --git a/platform/backend/app/core/scheduler.py b/platform/backend/app/core/scheduler.py index d7e4440..65d5a7a 100644 --- a/platform/backend/app/core/scheduler.py +++ b/platform/backend/app/core/scheduler.py @@ -25,6 +25,15 @@ class TaskScheduler: logger.warning("Scheduler already started") return # 使用 CronTrigger 设置每日固定时间点 + # 顺序: 搜索缓存(01:00)→采集(01:30)→趋势(03:00)→生成(03:30)→审查(04:30)→源优化(05:00)→指标(06:00) + self.scheduler.add_job( + self._run_refresh_search_cache, + CronTrigger(hour=1, minute=0), + id='scheduled_refresh_search_cache', + replace_existing=True, + max_instances=1, + coalesce=True + ) self.scheduler.add_job( self._run_collect, CronTrigger(hour=1, minute=30), @@ -62,14 +71,6 @@ class TaskScheduler: max_instances=1, coalesce=True ) - self.scheduler.add_job( - self._run_refresh_search_cache, - CronTrigger(hour=2, minute=30), - id='scheduled_refresh_search_cache', - replace_existing=True, - max_instances=1, - coalesce=True - ) self.scheduler.add_job( self._run_metrics_sync, CronTrigger(hour=6, minute=0), @@ -80,7 +81,7 @@ class TaskScheduler: ) self.scheduler.start() self._started = True - logger.info("Scheduler started: 01:30 collect, 02:30 refresh_search, 03:00 trends, 03:30 generate, 04:30 review, 05:00 optimize_sources, 06:00 metrics_sync") + logger.info("Scheduler started: 01:00 search_cache → 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() @@ -285,6 +286,36 @@ class TaskScheduler: if count: db.commit() logger.info("[Scheduled] Metrics sync completed: synced %d zhihu articles", count) + # 生成指标反馈:按 field 聚合表现,写入 metrics_feedback.json 供 collector 读取 + try: + import json as json_mod + from sqlalchemy import func as sql_func + feedback = db.query( + Topic.field, + sql_func.avg(ContentMetrics.likes).label("avg_likes"), + sql_func.avg(ContentMetrics.views).label("avg_views"), + sql_func.avg(ContentMetrics.comments).label("avg_comments"), + sql_func.count(ContentMetrics.id).label("article_count"), + ).join(ContentMetrics, ContentMetrics.topic_id == Topic.id + ).filter(Topic.field.isnot(None), Topic.field != "" + ).group_by(Topic.field).all() + if feedback: + scored = [] + for row in feedback: + score = (row.avg_likes or 0) + (row.avg_views or 0) * 0.01 + (row.avg_comments or 0) * 2 + scored.append((row.field, round(score, 1), int(row.article_count))) + scored.sort(key=lambda x: x[1], reverse=True) + feedback_data = { + "updated_at": datetime.now().isoformat(), + "top_domains": [(f, s) for f, s, _ in scored[:5]], + "detail": [{"field": f, "score": s, "articles": c} for f, s, c in scored], + } + feedback_file = Path(__file__).parent.parent.parent.parent / "automation" / "data" / "metrics_feedback.json" + feedback_file.parent.mkdir(parents=True, exist_ok=True) + feedback_file.write_text(json_mod.dumps(feedback_data, ensure_ascii=False, indent=2), encoding='utf-8') + logger.info("[Scheduled] Metrics feedback written: top domain %s (score %.1f)", scored[0][0], scored[0][1]) + except Exception as e_fb: + logger.warning("[Scheduled] Metrics feedback generation failed: %s", e_fb) else: logger.info("[Scheduled] Metrics sync: no zhihu articles to sync") db.close() diff --git a/scripts/collector.py b/scripts/collector.py index 7a0d5ac..2d41353 100644 --- a/scripts/collector.py +++ b/scripts/collector.py @@ -377,6 +377,61 @@ class SustainabilityCollector: logger.warning(f"web_search失败 {source.name}: {e}") return [] + def _get_trend_context(self) -> str: + """读取热点趋势数据和指标反馈, 返回markdown上下文""" + parts = [] + try: + from trends import load_trends + trends = load_trends() + if trends: + lines = ["## 当前热点趋势", ""] + for t in trends[:5]: + src = {"weibo": "🔥", "zhihu": "📖", "baidu": "🔍", "llm": "🤖"}.get(t.get("source", ""), "") + lines.append(f"- {src} **{t['topic']}**({t.get('platform','')}):{t.get('reason','')[:80]}") + kw = t.get("hot_keywords", []) + if kw: + lines.append(f" 搜索热词:{' '.join(kw[:3])}") + parts.append("\n".join(lines)) + except Exception: + pass + + metrics_file = DATA_DIR / "metrics_feedback.json" + if metrics_file.exists(): + try: + feedback = json.loads(metrics_file.read_text(encoding='utf-8')) + top_domains = feedback.get("top_domains", []) + if top_domains: + lines = ["## 历史表现反馈(高互动领域优先", ""] + for d, s in top_domains[:3]: + lines.append(f"- {d}:平均分 {s}") + parts.append("\n".join(lines)) + except Exception: + pass + + try: + from app.database import SessionLocal + from app.models import SystemConfig + db = SessionLocal() + try: + sc = db.query(SystemConfig).filter(SystemConfig.key == "collector_ai_advice").first() + if sc and sc.value: + advice = json.loads(sc.value) + summary = advice.get("summary", "") + new_cats = advice.get("suggested_new_categories", []) + if summary: + parts.append(f"## AI策略建议\n{summary}") + if new_cats: + suggested = [f"- {c['name']}({c.get('reason','')[:50]})" for c in new_cats[:2]] + parts.append("建议关注的新方向:\n" + "\n".join(suggested)) + except Exception: + pass + finally: + db.close() + except Exception: + pass + + return "\n\n".join(parts) + def _generate_topics_with_llm(self, cases: List[SustainabilityCase] = None, search_results: List[Dict] = None) -> List[SustainabilityTopic]: """用LLM基于采集数据生成选题(数据充分时精确生成,无数据时凭知识生成)""" try: @@ -402,11 +457,15 @@ class SustainabilityCollector: case_lines = [f"- {c.title[:40]}({c.category})" for c in cases[:5]] data_section += "\n采集案例:\n" + "\n".join(case_lines) + "\n" + trend_context = self._get_trend_context() + prompt = f"""你是一个内容策略师。基于以下信息,为「{target_category}」类别生成一个高质量选题。 {data_section if data_section else "(当前无实时采集数据,请基于你对中文互联网趋势的了解直接生成)"} {existing_hint} +{trend_context} + 输出一个选题,格式JSON: {{{{ "title": "标题(20字内,含核心关键词)",