From 8f2db92d9c441b8ee7966065630b8fc06df8bbd0 Mon Sep 17 00:00:00 2001 From: Simon Date: Wed, 12 Aug 2026 11:18:51 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20pipeline=20=E6=89=A7=E8=A1=8C=E6=97=B6?= =?UTF-8?q?=E6=98=BE=E6=80=A7=E8=BE=93=E5=87=BA=E5=BD=93=E5=89=8D=E9=98=B6?= =?UTF-8?q?=E6=AE=B5=E4=B8=8E=20AI=20=E5=A4=A7=E6=A8=A1=E5=9E=8B=E4=BF=A1?= =?UTF-8?q?=E6=81=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - run_pipeline 每步执行前打印阶段标题(序号/中文名/步骤名/日期) - llm/report/embedding 步骤打印 AI 大模型供应商与模型名 (provider/model 按 configs/llm_models.yaml 场景配置解析,缺 key 也可展示) --- scheduler/pipeline.py | 94 ++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 93 insertions(+), 1 deletion(-) diff --git a/scheduler/pipeline.py b/scheduler/pipeline.py index a16b1d8..993c54b 100644 --- a/scheduler/pipeline.py +++ b/scheduler/pipeline.py @@ -12,6 +12,7 @@ from __future__ import annotations import json +import os import subprocess import time from dataclasses import dataclass, field @@ -23,6 +24,27 @@ 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) @@ -122,6 +144,73 @@ def _resume_start_index( 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 步骤。 @@ -249,7 +338,10 @@ def run_pipeline( f" (跳过已成功 {names[:start_idx]})" if start_idx > 0 else "", ) - for name in names[start_idx:]: + 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)