打通流水线闭环: trends+metrics+optimize_sources → collector

三大数据流精准注入采集器:

1. 热点趋势(trends) -> collector
   _get_trend_context() 读取 trends.json
   热搜注入 LLM 选题 prompt

2. 历史表现(metrics) -> collector
   scheduler 同步指标后按field聚合写入 feedback.json
   collector 读取后高互动领域获优先级提升

3. AI策略(optimize_sources) -> collector
   从 DB SystemConfig 读取 AI 建议注入 prompt
This commit is contained in:
Yuzhiran Dev
2026-05-21 09:48:38 +08:00
parent 6b790e77f6
commit 56293a52a2
2 changed files with 99 additions and 9 deletions
+40 -9
View File
@@ -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()