初始化
This commit is contained in:
@@ -0,0 +1,231 @@
|
||||
"""全链路管道编排 (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:
|
||||
"""日报生成。"""
|
||||
started = datetime.now()
|
||||
try:
|
||||
from scheduler.reporter import generate_report
|
||||
path = generate_report()
|
||||
elapsed = (datetime.now() - started).total_seconds()
|
||||
ok = path is not None
|
||||
return StepResult(
|
||||
name="report", success=ok, elapsed_sec=elapsed,
|
||||
message=str(path) if path 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
|
||||
Reference in New Issue
Block a user