"""定时任务主流程(M7)。 编排 M1→M6 全链路,每一步调用已有脚本。 单步失败记录日志但不阻断后续(后续步骤可能使用旧缓存数据,降级继续)。 断点恢复: 每次运行把各步骤结果写入 data/pipeline/state.json(按日期隔离); run_pipeline(resume=True) 时跳过连续成功的步骤,从第一个失败/未执行 的步骤继续,实现 `pipeline --once --resume` 断点续跑。 """ from __future__ import annotations import json import os import subprocess import time from dataclasses import dataclass, field from datetime import date, datetime from pathlib import Path from loguru import logger # 断点状态文件(按日期隔离,记录每步骤结果) DEFAULT_STATE_PATH = Path("data/pipeline/state.json") # 步骤名 → 中文阶段名(显性输出用) _STAGE_LABELS: dict[str, str] = { "crawler": "M1 新闻抓取", "xwlb": "M1 新闻联播抓取", "extractor": "M2 正文提取", "dedup": "M3 新闻去重", "llm": "M4 LLM 事件抽取", "embedding": "M5 向量化嵌入", "qdrant": "M6 Qdrant 入库", "report": "日报生成", "cninfo_crawl": "cninfo 公告抓取", "cninfo_extract": "cninfo 正文提取", "cninfo_pdf": "cninfo PDF 补充", } # 使用对话大模型的步骤 → 对应 configs/llm_models.yaml 场景名 _AI_LLM_SCENES: dict[str, str] = { "llm": "event_extraction", "report": "daily_report", } # 步骤超时(秒) STEP_TIMEOUTS: dict[str, int] = { "crawler": 900, # M1 抓取(含 Playwright 浏览器,13 源约 8-12 min) "xwlb": 60, # M1 新闻联播 API(纯 HTTP,秒级) "extractor": 300, # M2 正文提取 "dedup": 120, # M3 去重 "llm": 900, # M4 LLM 事件抽取(API 调用,100 篇约 30s 但加限流余量) "embedding": 300, # M5 向量化 "qdrant": 300, # M6 入库(数据量大时需较长时间) "report": 30, # 日报生成+上传 "cninfo_crawl": 900, # cninfo watchlist URL 驱动(SPA 渲染,每只约 25s) } # 步骤对应的 uv run 命令(参数中 {date} 会被替换为实际日期) # crawler 不支持 --date,固定写当天目录; dedup 不加 --reset 以保持增量 STEP_COMMANDS: dict[str, list[str]] = { "crawler": ["uv", "run", "python", "-m", "scripts.run_crawler"], # 无 --date "xwlb": ["uv", "run", "python", "-m", "scripts.run_xwlb", "--date", "{date}"], # 处理日目录 "extractor": ["uv", "run", "python", "-m", "scripts.run_extractor", "--date", "{date}"], "dedup": ["uv", "run", "python", "-m", "scripts.run_dedup", "--date", "{date}"], "llm": ["uv", "run", "python", "-m", "scripts.run_event_extraction", "--date", "{date}"], "embedding": ["uv", "run", "python", "-m", "scripts.run_embedding", "--date", "{date}"], "qdrant": ["uv", "run", "python", "-m", "scripts.run_qdrant_ingest", "--date", "{date}"], "report": [], # 特殊步骤:仅在每天首次定时任务时追加 "cninfo_crawl": ["uv", "run", "a-share", "cninfo"], "cninfo_extract": ["uv", "run", "a-share", "extract", "--source", "cninfo", "--date", "{date}"], "cninfo_pdf": ["uv", "run", "a-share", "cninfo", "--enrich-pdf", "--pdf-limit", "100"], } @dataclass class StepResult: name: str success: bool elapsed_sec: float exit_code: int | None = None tail_msg: str = "" started_at: datetime | None = None @dataclass class PipelineResult: steps: list[StepResult] = field(default_factory=list) started_at: datetime | None = None finished_at: datetime | None = None @property def all_success(self) -> bool: return all(s.success for s in self.steps) def _load_pipeline_state(path: Path = DEFAULT_STATE_PATH) -> dict: """读取断点状态文件;不存在或损坏时返回空 dict。""" if not path.is_file(): return {} try: data = json.loads(path.read_text(encoding="utf-8")) except (json.JSONDecodeError, OSError) as e: logger.warning("pipeline 状态文件损坏,忽略: {} ({})", path, e) return {} return data if isinstance(data, dict) else {} def _save_pipeline_state(state: dict, path: Path = DEFAULT_STATE_PATH) -> None: """原子写状态文件(tmp + rename,避免中断写坏)。""" path.parent.mkdir(parents=True, exist_ok=True) tmp = path.with_suffix(".json.tmp") tmp.write_text(json.dumps(state, ensure_ascii=False, indent=2), encoding="utf-8") tmp.replace(path) def _update_step_state(state: dict, date_str: str, sr: StepResult) -> None: """把单步结果写入状态(ok/failed,含退出码与耗时)。""" day_state = state.setdefault(date_str, {}) day_state[sr.name] = { "status": "ok" if sr.success else "failed", "exit_code": sr.exit_code, "started_at": sr.started_at.isoformat() if sr.started_at else None, "elapsed_sec": round(sr.elapsed_sec, 1), } def _resume_start_index( names: list[str], date_str: str, state: dict, ) -> int: """计算断点续跑起始下标:跳过连续 ok 前缀,从首个失败/未记录步骤开始。 返回 0..len(names)-1;全部成功时返回 len(names)(表示无需续跑)。 """ day_state = state.get(date_str, {}) for i, name in enumerate(names): rec = day_state.get(name) if rec is None or rec.get("status") != "ok": return i return len(names) def _llm_scene_desc(scene: str) -> str | None: """解析某 LLM 场景的 provider/model(仅展示,不校验 API key)。""" try: from configs.loader import load_scene_config def _env(key: str) -> str | None: v = os.environ.get(key) return v.strip() if v else None sc = load_scene_config(scene) p = (sc.get("provider") or _env("LLM_PROVIDER") or "deepseek").lower() if p in ("qwen", "dashscope"): p = "qwen" model = sc.get("model") if not model: if p == "qwen": model = _env("QWEN_MODEL") or _env("LLM_MODEL") else: model = _env("DEEPSEEK_MODEL") or _env("LLM_MODEL") if not model: return None return f"provider={p}, model={model}" except Exception as e: # noqa: BLE001 - 配置缺失时降级展示 logger.debug("AI 描述解析失败(场景 {}): {}", scene, e) return None def _embedding_desc() -> str | None: """解析 embedding 场景的 provider/model(不实例化模型,避免加载本地权重)。""" try: from configs.loader import load_scene_config from embedding import resolve_provider_type sc = load_scene_config("embedding") pt = resolve_provider_type().value # dashscope | local-bge model = sc.get("model") if not model: if pt == "dashscope": model = os.environ.get("DASHSCOPE_EMBEDDING_MODEL") or "text-embedding-v3" else: model = os.environ.get("LOCAL_EMBEDDING_MODEL") or "BAAI/bge-m3" return f"provider={pt}, model={model}" except Exception as e: # noqa: BLE001 logger.debug("embedding 描述解析失败: {}", e) return None def _log_stage_header(name: str, date_str: str, index: int, total: int) -> None: """显性输出当前阶段(中文名 + 步骤名 + 序号 + 日期)。""" label = _STAGE_LABELS.get(name, name) logger.info("") logger.info("═" * 56) logger.info("阶段 {}/{}: {} [{}] 日期 {}", index, total, label, name, date_str) logger.info("═" * 56) def _log_ai_info(name: str) -> None: """本阶段用到 AI 大模型时,显性告知供应商与模型名。""" if name in _AI_LLM_SCENES: scene = _AI_LLM_SCENES[name] desc = _llm_scene_desc(scene) logger.info("🤖 本阶段使用 AI 大模型: {}", desc or "未配置(可能跳过或降级)") elif name == "embedding": desc = _embedding_desc() logger.info("🤖 本阶段使用 AI 嵌入模型: {}", desc or "未配置(可能降级)") def run_step(name: str, date_str: str) -> StepResult: """执行单个 pipeline 步骤。 参数: name: 步骤名(crawler/extractor/.../report) date_str: YYYYMMDD 日期字符串 返回: StepResult。 """ # 补跑保护:抓取类步骤只能产生当天数据(网站首页只含当前内容,历史文章已滚走), # 补跑历史日期时跳过抓取并告警,复用已有 data/raw/*/{date_str} 数据。 # 注意:仅保护 crawler; xwlb 走 API 支持任意历史日期,无需跳过。 if name == "crawler" and date_str != date.today().strftime("%Y%m%d"): started = datetime.now() logger.warning( "补跑模式:首页只含当天内容,无法补抓 {};跳过抓取,复用已有 data/raw/*/{}", date_str, date_str, ) return StepResult( name=name, success=True, elapsed_sec=0.0, tail_msg="补跑跳过(抓取只能产生当天数据)", started_at=started, ) # report 步骤:内部函数,不走子进程 # 日报按 date_str 日期生成: 新闻由 _collect_news_events 回溯过去 30 小时, # xwlb 由 _collect_xwlb 固定取前一日(已播出)联播。 if name == "report": started = datetime.now() try: from .reporter import generate_report # noqa: E402 report_date = date_str logger.info("日报: report_date={} (新闻 30h 回溯, xwlb 前一日)", report_date) path = generate_report(report_date, upload=True) elapsed = (datetime.now() - started).total_seconds() ok = path is not None return StepResult( name=name, success=ok, elapsed_sec=elapsed, tail_msg=str(path) if path else f"无数据 (report_date={report_date})", started_at=started, ) except Exception as e: # noqa: BLE001 elapsed = (datetime.now() - started).total_seconds() logger.exception("日报生成异常: {}", e) return StepResult(name=name, success=False, elapsed_sec=elapsed, tail_msg=str(e)[:200], started_at=started) cmd = STEP_COMMANDS.get(name) if cmd is None: return StepResult(name=name, success=False, elapsed_sec=0, tail_msg=f"未知步骤: {name}") full_cmd = [arg.replace("{date}", date_str) for arg in cmd] # 超时优先级: # 1. TIMEOUT_{NAME} 环境变量 (单步精确控制) # 2. PIPELINE_STEP_TIMEOUT 环境变量 (全局兜底, 覆盖硬编码) # 3. STEP_TIMEOUTS 硬编码字典 (代码内默认值) # 4. 1800s (最终兜底) import os specific_key = f"TIMEOUT_{name.upper()}" if specific_key in os.environ: timeout = int(os.environ[specific_key]) elif "PIPELINE_STEP_TIMEOUT" in os.environ: timeout = int(os.environ["PIPELINE_STEP_TIMEOUT"]) else: timeout = STEP_TIMEOUTS.get(name, 1800) started = datetime.now() logger.info("步骤 {} 开始: {}", name, " ".join(full_cmd)) try: # 不捕获输出,子进程日志直接流到终端(用户能看到每个源的抓取进度) proc = subprocess.run( full_cmd, timeout=timeout, ) elapsed = (datetime.now() - started).total_seconds() ok = proc.returncode == 0 # dedup 返回 1 是"重复率过高"(无新文章的正常场景) if name == "dedup" and proc.returncode == 1: ok = True logger.info("dedup 重复率超过阈值(无新数据场景,视为成功)") tail_msg = f"rc={proc.returncode}" if not ok else "" if ok: logger.info("步骤 {} 完成 ({}s) ✅", name, elapsed) else: logger.error("步骤 {} 失败 rc={} ({}s)", name, proc.returncode, elapsed) return StepResult( name=name, success=ok, elapsed_sec=elapsed, exit_code=proc.returncode, tail_msg=tail_msg, started_at=started, ) except subprocess.TimeoutExpired: elapsed = (datetime.now() - started).total_seconds() logger.error("步骤 {} 超时 (>{:.0f}s)", name, elapsed) return StepResult(name=name, success=False, elapsed_sec=elapsed, tail_msg="超时", started_at=started) except Exception as e: # noqa: BLE001 elapsed = (datetime.now() - started).total_seconds() logger.exception("步骤 {} 异常: {}", name, e) return StepResult(name=name, success=False, elapsed_sec=elapsed, tail_msg=str(e)[:200], started_at=started) def run_pipeline( date_str: str, *, steps: list[str] | None = None, resume: bool = False, state_path: Path = DEFAULT_STATE_PATH, ) -> PipelineResult: """串联执行全链路(M1→M6)。 参数: date_str: YYYYMMDD。 steps: 可选步骤列表,默认全部 6 步。 resume: True 时断点续跑——读取 data/pipeline/state.json 中该日期的 记录,跳过连续成功的步骤,从第一个失败/未执行步骤继续。 state_path: 断点状态文件路径(测试可注入)。 """ names = steps or [k for k in STEP_COMMANDS if k not in ("report", "cninfo_crawl", "cninfo_extract", "cninfo_pdf")] result = PipelineResult(started_at=datetime.now()) state = _load_pipeline_state(state_path) start_idx = 0 if resume: start_idx = _resume_start_index(names, date_str, state) if start_idx >= len(names): logger.info("resume: {} 的所有步骤均已完成,无需续跑", date_str) result.finished_at = datetime.now() return result logger.info( "resume: 从步骤 {} 继续{}", names[start_idx], f" (跳过已成功 {names[:start_idx]})" if start_idx > 0 else "", ) for i, name in enumerate(names[start_idx:], start=start_idx + 1): # 显性输出当前阶段 + AI 大模型信息 _log_stage_header(name, date_str, i, len(names)) _log_ai_info(name) sr = run_step(name, date_str) result.steps.append(sr) # 记录断点状态(无论成败,便于下次 resume) _update_step_state(state, date_str, sr) _save_pipeline_state(state, state_path) if not sr.success: logger.warning("步骤 {} 失败,后续步骤继续(可能降级)", name) # 步间留一点缓冲 time.sleep(0.5) result.finished_at = datetime.now() total = (result.finished_at - result.started_at).total_seconds() if result.started_at else 0 succ = sum(1 for s in result.steps if s.success) rate = succ / max(len(result.steps), 1) logger.info( "Pipeline 完成: {}/{} 步骤成功 ({:.0%}) 总耗时 {:.0f}s", succ, len(result.steps), rate, total, ) # 输出摘要 for s in result.steps: flag = "✅" if s.success else "❌" logger.info(" {} {} {}s", flag, s.name, s.elapsed_sec) return result