"""M7 定时任务调度入口。 支持两种模式: --once 立刻执行一次全链路 (默认) 启动 APScheduler,按 .env 中 SCHEDULE_TIMES 定时执行 用法: uv run python -m scripts.run_scheduler # 启动定时服务 uv run python -m scripts.run_scheduler --once # 立即执行一次 uv run python -m scripts.run_scheduler --once --date 20260616 uv run python -m scripts.run_scheduler --once --steps crawler,extractor """ from __future__ import annotations import argparse import json import signal import sys from pathlib import Path from typing import Any from loguru import logger from configs.runtime_env import ensure_env_loaded, env_raw, start_env_watcher from scheduler import DEFAULT_NEWS_STEPS, run_pipeline from scheduler.stock_reporter import generate_all_stock_reports from scheduler.timeutil import now as tz_now from scheduler.timeutil import schedule_tz, today_str def _setup_logger(level: str) -> None: logger.remove() logger.add( sys.stderr, level=level, format="{time:YYYY-MM-DD HH:mm:ss} | {level} | {name} | {message}", ) log_path = Path("logs") / "scheduler.log" log_path.parent.mkdir(parents=True, exist_ok=True) logger.add(log_path, level="DEBUG", rotation="10 MB", retention=5, encoding="utf-8") def _parse_schedule_times(raw: str) -> list[tuple[int, int]]: """解析 SCHEDULE_TIMES 环境变量。 格式: "07:00,12:00,18:00,22:00" 返回: [(7,0), (12,0), (18,0), (22,0)] """ out: list[tuple[int, int]] = [] for part in raw.split(","): part = part.strip() if not part: continue try: h, m = part.split(":") out.append((int(h), int(m))) except ValueError: logger.warning("SCHEDULE_TIMES 格式错误: {!r},跳过", part) return out def _parse_hhmm(raw: str, default: tuple[int, int]) -> tuple[int, int]: """解析 "HH:MM";非法或越界时记警告并返回 default。""" parts = raw.split(":") try: hour = int(parts[0]) minute = int(parts[1]) if len(parts) > 1 else 0 except (ValueError, IndexError): logger.warning("时间格式错误: {!r},改用默认 {:02d}:{:02d}", raw, *default) return default if not (0 <= hour <= 23 and 0 <= minute <= 59): logger.warning("时间越界: {!r},改用默认 {:02d}:{:02d}", raw, *default) return default return hour, minute def _scheduled_pipeline(steps: list[str] | None = None) -> None: """定时触发的全链路:每次触发时按调度时区重新计算日期。""" run_pipeline(today_str(), steps=steps) def _desired_jobs() -> dict[str, dict[str, Any]]: """按当前环境变量算出「期望的」定时任务集合。 每次调用都通过 env_get 读取,所以改 .env 后无需重启即可反映到 _sync_jobs。 """ jobs: dict[str, dict[str, Any]] = {} times = _parse_schedule_times( env_raw("SCHEDULE_TIMES", "07:00,12:00,18:00,22:00") or "" ) if times: first = min(times) for hour, minute in times: with_report = (hour, minute) == first jobs[f"pipeline_{hour:02d}{minute:02d}"] = { "kind": "pipeline", "hour": hour, "minute": minute, "steps": list(DEFAULT_NEWS_STEPS) + (["report"] if with_report else []), "name": f"全链路 {'+日报' if with_report else ''} {hour:02d}:{minute:02d}", } cninfo_h, cninfo_m = _parse_hhmm( env_raw("CNINFO_SCHEDULE_TIME", "06:30") or "06:30", (6, 30) ) jobs["pipeline_cninfo"] = { "kind": "pipeline", "hour": cninfo_h, "minute": cninfo_m, "steps": ["cninfo_crawl", "cninfo_extract", "cninfo_pdf", "dedup", "llm", "embedding", "qdrant"], "name": f"cninfo 公告管道 {cninfo_h:02d}:{cninfo_m:02d}", } # STOCK_REPORT_TIME 缺省 07:30;显式留空表示禁用 stock_raw = env_raw("STOCK_REPORT_TIME") if stock_raw is None: stock_raw = "07:30" if stock_raw: stock_h, stock_m = _parse_hhmm(stock_raw, (7, 30)) jobs["stock_report"] = { "kind": "stock", "hour": stock_h, "minute": stock_m, "name": f"个股日报 {stock_h:02d}:{stock_m:02d}", } return jobs _jobs_sig: str | None = None def _sync_jobs(scheduler: Any) -> bool: """把「期望任务」同步到 APScheduler;只在配置变化时增删。返回是否变更。 以 ``_`` 开头的内部任务(如配置热同步自身)不参与增删。 """ global _jobs_sig from apscheduler.triggers.cron import CronTrigger # noqa: E402 desired = _desired_jobs() if not any(jid.startswith("pipeline_") for jid in desired): logger.error("SCHEDULE_TIMES 为空或全部非法,保留现有定时任务不改动") return False sig = json.dumps(desired, sort_keys=True, ensure_ascii=False) if sig == _jobs_sig: return False for jid, spec in desired.items(): trigger = CronTrigger( hour=spec["hour"], minute=spec["minute"], timezone=str(schedule_tz()) ) if spec["kind"] == "pipeline": scheduler.add_job( _scheduled_pipeline, trigger=trigger, kwargs={"steps": spec["steps"]}, id=jid, name=spec["name"], replace_existing=True, ) else: scheduler.add_job( generate_all_stock_reports, trigger=trigger, id=jid, name=spec["name"], replace_existing=True, ) for job in scheduler.get_jobs(): if job.id.startswith("_") or job.id in desired: continue scheduler.remove_job(job.id) _jobs_sig = sig logger.info( "定时任务已同步: {}", ", ".join(f"{s['name']}" for s in desired.values()), ) return True def _once(args: argparse.Namespace) -> int: """单次执行模式。 默认全量执行;--resume 时断点续跑(跳过连续成功步骤,从失败/未执行步骤继续)。 """ steps = None if args.steps: steps = [s.strip() for s in args.steps.split(",")] if args.resume and args.steps: logger.error("--resume 与 --steps 不能同时使用(断点续跑针对全链路)") return 2 run_pipeline(args.date or today_str(), steps=steps, resume=args.resume) return 0 def _daemon(args: argparse.Namespace) -> int: """守护进程模式(APScheduler)。""" from apscheduler.schedulers.background import BackgroundScheduler # noqa: E402 times = _parse_schedule_times( env_raw("SCHEDULE_TIMES", "07:00,12:00,18:00,22:00") or "" ) if not times: logger.error("SCHEDULE_TIMES 为空或全部非法,无法启动定时任务") return 2 sorted_times = sorted(times) scheduler = BackgroundScheduler() # 首次注册;此后由 _config_watch 每 30s 热同步,改 .env 无需重启 _sync_jobs(scheduler) # 优雅退出 def _shutdown(signum: int, frame: Any) -> None: logger.info("收到信号 {}, 关闭调度器...", signum) scheduler.shutdown(wait=False) raise SystemExit(0) signal.signal(signal.SIGINT, _shutdown) signal.signal(signal.SIGTERM, _shutdown) scheduler.start() # 配置热同步:每 30s 重新计算 SCHEDULE_TIMES / CNINFO_SCHEDULE_TIME / STOCK_REPORT_TIME scheduler.add_job( _sync_jobs, "interval", seconds=30, args=[scheduler], id="_config_watch", name="配置热同步", replace_existing=True, ) # .env 热加载监听:provider / model / key / base_url 等改动无需重启 start_env_watcher() logger.info("调度器已启动,等待触发... (按 Ctrl+C 退出;改 .env 无需重启)") # 启动时检查是否有因重启/宕机错过的定时任务,30 分钟内补跑 now = tz_now() for hour, minute in times: scheduled = now.replace(hour=hour, minute=minute, second=0, microsecond=0) missed_minutes = (now - scheduled).total_seconds() / 60 if 0 < missed_minutes < 30: logger.warning( "检测到错过的定时任务 {:02d}:{:02d} ({} 分钟前),立即补跑一次", hour, minute, int(missed_minutes), ) # P1-2:补跑只跑新闻链路(DEFAULT_NEWS_STEPS 不含 report 与 cninfo 三步), # cninfo 公告管道由其自身定时任务负责,避免重复执行整套公告管道。 steps = list(DEFAULT_NEWS_STEPS) if (hour, minute) == sorted_times[0]: steps.append("report") run_pipeline(today_str(), steps=steps) import contextlib with contextlib.suppress(SystemExit, KeyboardInterrupt): # 保持主线程存活,直到收到退出信号 signal.pause() return 0 def main() -> int: parser = argparse.ArgumentParser(description="A 股新闻定时任务 (M7)") parser.add_argument("--once", action="store_true", help="立即执行一次全链路") parser.add_argument( "--date", default=None, help="日期 YYYYMMDD (仅 --once 模式,默认按调度时区取今天)", ) parser.add_argument("--steps", default=None, help="仅执行指定步骤,逗号分隔 (如 crawler,extractor)") parser.add_argument( "--resume", action="store_true", help="断点续跑(仅 --once):跳过连续成功步骤,从上次失败/未执行步骤继续", ) parser.add_argument("--log-level", default="INFO") args = parser.parse_args() _setup_logger(args.log_level) # .env 热加载(改文件后常驻进程无需重启;--once 也会即时读取最新配置) ensure_env_loaded() if args.once: return _once(args) return _daemon(args) if __name__ == "__main__": raise SystemExit(main())