Files
news/scheduler/pipeline.py
T
simon 8fa27ad65b fix: 修复 dedup 高重复率返回码被掩蔽的问题 (P1-1)
- run_dedup: 生产模式重复率仅作 WARNING 告警,不影响退出码(执行成功即 0)
- run_dedup: 新增 --strict 验收模式(重复率 > 5% 返回 1),保留 M3 验收门槛
- run_dedup: 新增统计快照 data/deduped/{date}/stats.json(原子写)
- run_dedup: 顺带修复空日场景 sources.json 写入 FileNotFoundError
- pipeline: 移除 dedup rc=1 特判,恢复'非 0 即失败'统一语义
- 新增 tests/test_run_dedup.py 7 个测试;全量 266 passed
2026-08-22 19:50:14 +08:00

379 lines
15 KiB
Python

"""定时任务主流程(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
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