#!/bin/bash # ============================================= # 国内服务器:全链路管道 M2 → M3 → M4 → M5 → M6 → 日报 # ============================================= # 用法: # ./scripts/pipeline.sh # 全新执行(步骤级不跳过) # ./scripts/pipeline.sh --resume # 从中断处继续(跳过已完成步骤) # 前提:domestic_sync.sh 已完成 # ============================================= # 各步骤"跳过已处理文件"(文件级增量,由各 Python 模块自带): # M2 正文提取: data/processed/{source}/{date}/{url_hash}.json 存在则跳过 # M3 去重: 指纹库 SQLite(data/dedup_fingerprints.sqlite3)判重,重复自动丢弃 # M4 翻译: data/events/{date}/{url_hash}.json 存在则跳过 # M5 向量: data/embeddings/{date}/index.json 记录已向量化,跳过 # M6 入库: Qdrant upsert 幂等(点 id = url_hash,重复写入覆盖,无重复) # 日报: MySQL 唯一键 (report_date, intl, "") 幂等覆盖 # ============================================= set -euo pipefail SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" PROJECT_DIR="$(dirname "$SCRIPT_DIR")" cd "$PROJECT_DIR" # ── 参数解析 ── RESUME=0 for arg in "$@"; do case "$arg" in --resume) RESUME=1 ;; *) echo "未知参数: ${arg}(支持 --resume)" >&2; exit 1 ;; esac done source "$SCRIPT_DIR/_step_state.sh" LOG() { echo "[$(date '+%Y-%m-%d %H:%M:%S')] $*"; } # ── 步骤日志:全部输出实时显示到终端,同时完整写入日志文件 ── LOG_DIR="logs" mkdir -p "$LOG_DIR" STEP_LOG="$LOG_DIR/pipeline_$(date +%Y%m%d_%H%M%S).log" : > "$STEP_LOG" LOG "══════ 全链路管道开始(resume=${RESUME})═══════" LOG "步骤日志: ${STEP_LOG}(终端实时显示完整进度)" # ── AI 模型信息读取(供阶段 banner 显性展示供应商/模型)── ai_model_info() { # $1: 场景名(translation / daily_report)或 "embedding" # 输出格式: "供应商 / 模型" .venv/bin/python3 -c " import sys, yaml scene = sys.argv[1] raw = yaml.safe_load(open('configs/system.yaml', encoding='utf-8')) if scene == 'embedding': cfg = raw.get('embedding', {}) p = cfg.get('provider', '?') m = cfg.get('dashscope_model', '?') else: base = raw.get('llm', {}) sc = (raw.get('llm_scenes', {}) or {}).get(scene, {}) merged = {**base, **sc} p = merged.get('provider', '?') m = merged.get('model') or merged.get(p + '_model', '?') print(f'{p} / {m}') " "$1" } # 加载 .env export $(grep -v '^#' .env | grep -v '^$' | xargs 2>/dev/null || true) # ── M2: 正文提取 ── LOG "━━━ M2 正文提取 ━━━" step_run M2_extract .venv/bin/python3 -c " from extractor.pipeline import process_all_sources stats = process_all_sources() print(f'M2: {stats[\"total_articles\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s') " # ── M3: 去重 ── LOG "━━━ M3 三层去重 ━━━" step_run M3_dedup .venv/bin/python3 -c " from dedup.pipeline import dedup_all_sources stats = dedup_all_sources() print(f'M3: 唯一 {stats[\"unique\"]}/重复 {stats[\"duplicate\"]}, {stats[\"elapsed_sec\"]:.0f}s') " # ── M4: 翻译+事件(AI 大模型)── LOG "━━━ M4 翻译+事件抽取(AI 大模型: $(ai_model_info translation))━━━" step_run M4_translate .venv/bin/python3 -c " from llm.pipeline import translate_all_deduped stats = translate_all_deduped() print(f'M4: {stats[\"success\"]}/{stats[\"total\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s') " # ── M5: 向量生成(AI 大模型)── LOG "━━━ M5 向量生成(AI 大模型: $(ai_model_info embedding))━━━" step_run M5_embed .venv/bin/python3 -c " from embedding.pipeline import embed_all_events stats = embed_all_events() print(f'M5: {stats[\"success\"]}/{stats[\"total\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s') " # ── M6: Qdrant 入库 ── LOG "━━━ M6 Qdrant 入库 ━━━" step_run M6_index .venv/bin/python3 -c " from vectorstore.pipeline import ingest_all_embeddings stats = ingest_all_embeddings() print(f'M6: {stats[\"ingested\"]}/{stats[\"total\"]} 条, {stats[\"elapsed_sec\"]:.0f}s') " # ── 日报(AI 大模型)── LOG "━━━ 日报生成(AI 大模型: $(ai_model_info daily_report))━━━" step_run report .venv/bin/python3 -c " from scheduler.reporter import generate_report report_id = generate_report() print(f'日报: report_id={report_id}' if report_id is not None else '日报: 无数据/失败') " LOG "══════ 全链路管道完成 ✅ ═══════"