fef435cc78
- 审查:移除 manual_review,改为迭代LLM修复(最多3次),合规分回写Topic - 调度:scheduler 新增话题采集定时任务 scheduled_collect (01:30) - 提示词:全链路8文件≈24个提示词升级,增强SEO/平台推荐/真人感
119 lines
4.2 KiB
Python
119 lines
4.2 KiB
Python
"""
|
|
定时任务调度器
|
|
基于 APScheduler,支持在 FastAPI 生命周期内运行定时任务
|
|
"""
|
|
import os
|
|
import logging
|
|
from datetime import datetime
|
|
from apscheduler.schedulers.background import BackgroundScheduler
|
|
from apscheduler.triggers.cron import CronTrigger
|
|
from .generator import run_creator
|
|
from .optimizer import run_optimizer
|
|
from .sync import sync_all_topics
|
|
from .collector import run_collector
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
class TaskScheduler:
|
|
def __init__(self):
|
|
self.scheduler = BackgroundScheduler()
|
|
self._started = False
|
|
|
|
def start(self):
|
|
if self._started:
|
|
logger.warning("Scheduler already started")
|
|
return
|
|
# 使用 CronTrigger 设置每日固定时间点
|
|
self.scheduler.add_job(
|
|
self._run_collect,
|
|
CronTrigger(hour=1, minute=30),
|
|
id='scheduled_collect',
|
|
replace_existing=True,
|
|
max_instances=1,
|
|
coalesce=True
|
|
)
|
|
self.scheduler.add_job(
|
|
self._run_sync,
|
|
CronTrigger(hour=2, minute=30),
|
|
id='scheduled_sync',
|
|
replace_existing=True,
|
|
max_instances=1,
|
|
coalesce=True
|
|
)
|
|
self.scheduler.add_job(
|
|
self._run_generate,
|
|
CronTrigger(hour=3, minute=30),
|
|
id='scheduled_generate',
|
|
replace_existing=True,
|
|
max_instances=1,
|
|
coalesce=True
|
|
)
|
|
self.scheduler.add_job(
|
|
self._run_optimize,
|
|
CronTrigger(hour=4, minute=30),
|
|
id='scheduled_optimize',
|
|
replace_existing=True,
|
|
max_instances=1,
|
|
coalesce=True
|
|
)
|
|
self.scheduler.start()
|
|
self._started = True
|
|
logger.info("Scheduler started with daily cron triggers (01:30 collect, 02:30 sync, 03:30 generate, 04:30 optimize)")
|
|
def shutdown(self):
|
|
if self.scheduler.running:
|
|
self.scheduler.shutdown()
|
|
logger.info("Scheduler shut down")
|
|
|
|
def _run_generate(self):
|
|
try:
|
|
logger.info("[Scheduled] Starting content generation...")
|
|
result = run_creator()
|
|
logger.info("[Scheduled] Generation completed: %s", result)
|
|
created_id = result.get("topic_id") if isinstance(result, dict) else None
|
|
if created_id:
|
|
logger.info("[Scheduled] Running compliance review on %s...", created_id)
|
|
review_result = run_optimizer([created_id])
|
|
if review_result.get("ok"):
|
|
logger.info("[Scheduled] Review completed for %s", created_id)
|
|
else:
|
|
logger.warning("[Scheduled] Review failed: %s", review_result.get("error"))
|
|
except Exception as e:
|
|
logger.exception("[Scheduled] Generation pipeline failed: %s", e)
|
|
|
|
def _run_optimize(self):
|
|
try:
|
|
logger.info("[Scheduled] Starting compliance optimization...")
|
|
result = run_optimizer()
|
|
logger.info("[Scheduled] Optimization completed: %s", result)
|
|
except Exception as e:
|
|
logger.exception("[Scheduled] Optimization failed: %s", e)
|
|
|
|
def _run_collect(self):
|
|
try:
|
|
logger.info("[Scheduled] Starting topic collection...")
|
|
result = run_collector()
|
|
logger.info("[Scheduled] Collection completed: %s", result.get("output", "")[-200:])
|
|
except Exception as e:
|
|
logger.exception("[Scheduled] Collection failed: %s", e)
|
|
|
|
def _run_sync(self):
|
|
try:
|
|
logger.info("[Scheduled] Starting data sync...")
|
|
sync_all_topics()
|
|
logger.info("[Scheduled] Sync completed")
|
|
except Exception as e:
|
|
logger.exception("[Scheduled] Sync failed: %s", e)
|
|
|
|
def get_jobs(self):
|
|
"""返回当前所有定时任务的状态"""
|
|
jobs = []
|
|
for job in self.scheduler.get_jobs():
|
|
jobs.append({
|
|
"id": job.id,
|
|
"next_run_time": job.next_run_time.isoformat() if job.next_run_time else None,
|
|
"trigger": str(job.trigger),
|
|
})
|
|
return jobs
|
|
|
|
# 全局单例
|
|
scheduler = TaskScheduler() |