Files
intl_news/scheduler/pipeline.py
T
simon 3b44f64f66 refactor: 清理历史 AI Agent 文档残留 + 重构 docs/ + Pipeline 健壮性修复
- docs: 删除 CLAUDE.md / continuation.md / english-news-plan.md 及旧版 intlnews_usage.*,
        统一迁移到 docs/{README,architecture,quickstart,usage,pipeline,configuration,deployment,development,faq}.md
- README: 精简为仓库入口,指向 docs/
- configs/sources.yaml: 更新注释指向新文档
- .env.example: 修正 DashScope Embedding 端点说明

Pipeline 修复:
- dedup/llm/embedding/vectorstore/reporter: 过滤 M2 no_content / 空正文,避免污染下游与 Qdrant
- dedup/pipeline: 改为先写唯一文件再写指纹,避免崩溃导致文章永久丢失
- crawler/orchestrator: sources_crawled 改为“尝试数”,成功数 = crawled - failed
- crawler/storage: write_index_jsonl 从文章路径推断日期,修复跨天/测试路径问题
- scheduler/pipeline: STEP_TIMEOUTS 实际生效(SIGALRM)
- scheduler/reporter: emb_count 排除 index.json;日报跳过无原文事件
- vectorstore/pipeline: payload 增加 source_ids;--recreate --all 时空日期也重建 collection
- app/cli: extract/dedup/translate/embed/index/pipeline 支持 --date;embed/index 支持 --all;crawl 全源失败返回非零
- scripts: domestic_full/crawl_8g/crawl_2g/pipeline 安全加载 .env;M1 全失败不标记且最终退出码=1
2026-08-22 20:47:53 +08:00

277 lines
9.0 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""全链路管道编排 (M7)。
串联 M2 → M3 → M4 → M5 → M6,每步失败记录日志但不阻断后续。
"""
import logging
import signal
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,
}
class _StepTimeout(Exception):
"""步骤超时专用异常,避免与业务 TimeoutError 混淆。"""
def _run_step_with_timeout(func, name: str, date_str: str, timeout: int | None) -> "StepResult":
"""在支持 SIGALRM 的主线程中为单步执行添加超时保护。"""
if timeout is None:
return func(date_str)
started = datetime.now()
if not hasattr(signal, "SIGALRM"):
return func(date_str)
def _handler(signum, frame): # noqa: ARG001
raise _StepTimeout(f"step {name} timed out after {timeout}s")
try:
old_handler = signal.getsignal(signal.SIGALRM)
except (ValueError, OSError):
# 非主线程无法设置信号处理器,直接不启用超时
return func(date_str)
signal.signal(signal.SIGALRM, _handler)
signal.setitimer(signal.ITIMER_REAL, timeout)
try:
return func(date_str)
except _StepTimeout:
elapsed = (datetime.now() - started).total_seconds()
return StepResult(
name=name,
success=False,
elapsed_sec=elapsed,
message=f"超时(>{timeout}s)",
started_at=started,
)
finally:
signal.setitimer(signal.ITIMER_REAL, 0)
signal.signal(signal.SIGALRM, old_handler)
@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 = _run_step_with_timeout(
func, name, date_str, STEP_TIMEOUTS.get(name)
)
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