Files
intl_news/scripts/pipeline.sh
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

124 lines
4.6 KiB
Bash
Executable File
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.
#!/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
if [ -f .env ]; then
set -a
. ./.env
set +a
fi
# ── 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 "══════ 全链路管道完成 ✅ ═══════"