"""全链路管道编排 (M7)。 串联 M2 → M3 → M4 → M5 → M6,每步失败记录日志但不阻断后续。 """ import logging import time from dataclasses import dataclass, field from datetime import datetime from dotenv import load_dotenv load_dotenv() # 确保 .env 中的 API Key 被加载到 os.environ # 抑制第三方库噪音日志 logging.getLogger("httpx").setLevel(logging.WARNING) logging.getLogger("httpcore").setLevel(logging.WARNING) logging.getLogger("openai").setLevel(logging.WARNING) logger = logging.getLogger(__name__) # 步骤默认超时(秒) STEP_TIMEOUTS: dict[str, int] = { "extract": 300, "dedup": 120, "translate": 900, "embed": 300, "index": 300, "report": 60, } @dataclass class StepResult: """单步执行结果。""" name: str success: bool elapsed_sec: float message: 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) @property def success_count(self) -> int: return sum(1 for s in self.steps if s.success) def run_step_extract(date_str: str) -> StepResult: """M2: 正文提取。""" started = datetime.now() try: from extractor.pipeline import process_all_sources stats = process_all_sources(date_str=date_str) elapsed = (datetime.now() - started).total_seconds() return StepResult( name="extract", success=True, elapsed_sec=elapsed, message=f"{stats['total_articles']} 篇", started_at=started, ) except Exception as e: elapsed = (datetime.now() - started).total_seconds() logger.exception("M2 正文提取失败") return StepResult(name="extract", success=False, elapsed_sec=elapsed, message=str(e)[:200], started_at=started) def run_step_dedup(date_str: str) -> StepResult: """M3: 三层去重。""" started = datetime.now() try: from dedup.pipeline import dedup_all_sources stats = dedup_all_sources(date_str=date_str) elapsed = (datetime.now() - started).total_seconds() return StepResult( name="dedup", success=True, elapsed_sec=elapsed, message=f"唯一 {stats['unique']}/重复 {stats['duplicate']}", started_at=started, ) except Exception as e: elapsed = (datetime.now() - started).total_seconds() logger.exception("M3 去重失败") return StepResult(name="dedup", success=False, elapsed_sec=elapsed, message=str(e)[:200], started_at=started) def run_step_translate(date_str: str) -> StepResult: """M4: 翻译 + 事件抽取。""" started = datetime.now() try: from llm.pipeline import translate_all_deduped stats = translate_all_deduped(date_str=date_str) elapsed = (datetime.now() - started).total_seconds() return StepResult( name="translate", success=True, elapsed_sec=elapsed, message=f"{stats['success']}/{stats['total']} 篇 ({stats.get('provider','')})", started_at=started, ) except Exception as e: elapsed = (datetime.now() - started).total_seconds() logger.exception("M4 翻译失败") return StepResult(name="translate", success=False, elapsed_sec=elapsed, message=str(e)[:200], started_at=started) def run_step_embed(date_str: str) -> StepResult: """M5: 向量生成。""" started = datetime.now() try: from embedding.pipeline import embed_all_events stats = embed_all_events(date_str=date_str) elapsed = (datetime.now() - started).total_seconds() return StepResult( name="embed", success=True, elapsed_sec=elapsed, message=f"{stats['success']}/{stats['total']} 篇", started_at=started, ) except Exception as e: elapsed = (datetime.now() - started).total_seconds() logger.exception("M5 向量生成失败") return StepResult(name="embed", success=False, elapsed_sec=elapsed, message=str(e)[:200], started_at=started) def run_step_index(date_str: str) -> StepResult: """M6: Qdrant 入库。""" started = datetime.now() try: from vectorstore.pipeline import ingest_all_embeddings stats = ingest_all_embeddings(date_str=date_str) elapsed = (datetime.now() - started).total_seconds() return StepResult( name="index", success=True, elapsed_sec=elapsed, message=f"{stats['ingested']}/{stats['total']} 条", started_at=started, ) except Exception as e: elapsed = (datetime.now() - started).total_seconds() logger.exception("M6 入库失败") return StepResult(name="index", success=False, elapsed_sec=elapsed, message=str(e)[:200], started_at=started) def run_step_report(date_str: str) -> StepResult: """日报生成(M9:结构化入库,不再产出 HTML)。""" started = datetime.now() try: from scheduler.reporter import generate_report report_id = generate_report() elapsed = (datetime.now() - started).total_seconds() ok = report_id is not None return StepResult( name="report", success=ok, elapsed_sec=elapsed, message=f"report_id={report_id}" if report_id is not None else "无数据", started_at=started, ) except Exception as e: elapsed = (datetime.now() - started).total_seconds() logger.exception("日报生成异常") return StepResult(name="report", success=False, elapsed_sec=elapsed, message=str(e)[:200], started_at=started) def run_pipeline( date_str: str, *, steps: list[str] | None = None, skip_report: bool = False, ) -> PipelineResult: """串联执行全链路 M2→M6(+ 可选日报)。 Args: date_str: YYYYMMDD 日期 steps: 可选步骤列表,默认全部 skip_report: 是否跳过日报生成 Returns: PipelineResult """ if steps is None: steps = ["extract", "dedup", "translate", "embed", "index"] if not skip_report: steps.append("report") step_funcs = { "extract": run_step_extract, "dedup": run_step_dedup, "translate": run_step_translate, "embed": run_step_embed, "index": run_step_index, "report": run_step_report, } result = PipelineResult(started_at=datetime.now()) for name in steps: func = step_funcs.get(name) if func is None: logger.warning("未知步骤: %s,跳过", name) result.steps.append(StepResult(name=name, success=False, elapsed_sec=0, message=f"未知步骤: {name}")) continue logger.info("── 步骤 %s 开始 ──", name) sr = func(date_str) result.steps.append(sr) flag = "✅" if sr.success else "❌" logger.info("── 步骤 %s %s (%.1fs) %s", name, flag, sr.elapsed_sec, sr.message) if not sr.success: logger.warning("步骤 %s 失败,后续步骤继续", 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 logger.info( "Pipeline 完成: %d/%d 步骤成功,总耗时 %.0fs", result.success_count, len(result.steps), total, ) return result