Files
news/scheduler/pipeline.py
T
simon 7b33182f67 fix: 修复补跑误含 cninfo 三步与时区不一致 (P1-2/P1-3)
- P1-2: 提取 DEFAULT_NEWS_STEPS(全链路去 report+cninfo),定时/补跑/默认三处统一;
  启动补跑不再误执行 cninfo 公告管道
- P1-3: 新增 scheduler/timeutil.py(调度时区统一入口,SCHEDULE_TZ 可覆盖,
  默认 Asia/Shanghai);run_scheduler 定时/补跑/--once、crawler 补跑保护、
  reporter 兜底日期全部改用调度时区,避免系统时区非上海时日期错位
- 新增 6 个测试;全量 282 passed
2026-08-23 05:40:52 +08:00

390 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 datetime
from pathlib import Path
from loguru import logger
from .timeutil import today_str
# 断点状态文件(按日期隔离,记录每步骤结果)
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"],
}
# 新闻链路默认步骤:全链路除去 report(单独追加)与 cninfo 独立管道(P1-2:
# 补跑/定时任务若误含 cninfo 三步,会重复执行整套公告管道,且 cninfo_crawl
# 无 --date 参数,补跑历史日期时会静默空转)。
DEFAULT_NEWS_STEPS: list[str] = [
k for k in STEP_COMMANDS
if k not in ("report", "cninfo_crawl", "cninfo_extract", "cninfo_pdf")
]
@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 支持任意历史日期,无需跳过。
# "今天"按调度时区判定(P1-3),避免系统时区与 cron 时区不一致时错判。
if name == "crawler" and date_str != today_str():
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 list(DEFAULT_NEWS_STEPS)
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