- 新增 report_db/ 包(models/schema/db,复用 news 项目实现,幂等 upsert) - reporter.py: _build_report_data + generate_report 写库返回 report_id - pipeline/cli 适配 report_id 返回值;HTML 渲染/上传保留 deprecated - 新增 tests/test_report_db.py;.env.example 增加 NEWS_DB_* 配置 - 已部署 pi5 并验证:真实生成 report_id=184(12 事件)+ 幂等覆盖
233 lines
7.7 KiB
Python
233 lines
7.7 KiB
Python
"""全链路管道编排 (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
|