feat: 全流程脚本 --resume 断点续跑 + 各步骤文件级增量说明
- 新增 scripts/_step_state.sh:步骤状态库(data/run_state/{date}.state,按天换新)
- pipeline.sh / domestic_full.sh 支持 --resume,按步骤跳过已完成项
- step_run 显式检查退出码(规避 bash set -e 条件上下文陷阱),失败步骤不标记
- M1 保持部分源失败容忍(step_should_run/step_mark 手动控制)
- 确认并文档化各步骤文件级增量:M2/M4/M5 按文件存在跳过、M3 指纹库、M6 upsert 幂等
- 已同步 pi5 并验证(--badarg exit=1、step_run 机制正常)
This commit is contained in:
@@ -0,0 +1,60 @@
|
||||
#!/bin/bash
|
||||
# =============================================
|
||||
# 步骤状态管理(--resume 中断恢复用)
|
||||
# =============================================
|
||||
# 用法(在脚本中 source):
|
||||
# source "$SCRIPT_DIR/_step_state.sh"
|
||||
# RESUME=0|1 # 由入口脚本参数解析后设置
|
||||
# step_run <步骤名> <cmd...> # 执行并在成功后标记;resume 时已标记则跳过
|
||||
# step_should_run <名> # 返回 0 = 需要运行(供需容忍部分失败的步骤手动控制)
|
||||
# step_mark <名> # 手动标记完成
|
||||
# =============================================
|
||||
# 状态文件: data/run_state/{YYYYMMDD}.state(按本地日期换新,跨天自动失效)
|
||||
# 说明: 各步骤内部的"跳过已处理文件"由各 Python 模块自带(见 pipeline.sh 头注释),
|
||||
# 本文件仅做"步骤级"断点记录。
|
||||
# =============================================
|
||||
|
||||
STEP_STATE_DIR="data/run_state"
|
||||
|
||||
step_state_file() {
|
||||
echo "${STEP_STATE_DIR}/$(date +%Y%m%d).state"
|
||||
}
|
||||
|
||||
step_done() {
|
||||
# $1: 步骤名;已标记返回 0
|
||||
[ -f "$(step_state_file)" ] && grep -qx "$1" "$(step_state_file)" >/dev/null 2>&1
|
||||
}
|
||||
|
||||
step_mark() {
|
||||
# $1: 步骤名;记录为已完成(幂等追加)
|
||||
mkdir -p "${STEP_STATE_DIR}"
|
||||
echo "$1" >> "$(step_state_file)"
|
||||
}
|
||||
|
||||
step_should_run() {
|
||||
# $1: 步骤名;resume 模式下已标记 → 返回 1(跳过),否则返回 0
|
||||
if [ "${RESUME:-0}" = "1" ] && step_done "$1"; then
|
||||
LOG "⏭️ 跳过已完成的步骤: $1(--resume)"
|
||||
return 1
|
||||
fi
|
||||
return 0
|
||||
}
|
||||
|
||||
step_run() {
|
||||
# $1: 步骤名,$2...: 要执行的命令(成功后标记完成)
|
||||
# 显式检查退出码:避免 bash set -e 在条件上下文(&&/||)调用函数时失效,
|
||||
# 导致失败步骤被误标记为完成。
|
||||
local name="$1"
|
||||
shift
|
||||
if ! step_should_run "$name"; then
|
||||
return 0
|
||||
fi
|
||||
LOG "▶ $name"
|
||||
"$@"
|
||||
local rc=$?
|
||||
if [ "$rc" -ne 0 ]; then
|
||||
LOG "✗ $name 失败(exit=$rc),未标记完成;可 --resume 重试该步骤"
|
||||
return "$rc"
|
||||
fi
|
||||
step_mark "$name"
|
||||
}
|
||||
@@ -6,25 +6,50 @@
|
||||
# =============================================
|
||||
# 每天 06:00 / 12:00 / 18:00 / 22:00 各执行一次
|
||||
# =============================================
|
||||
# 用法:
|
||||
# ./scripts/domestic_full.sh # 全新执行
|
||||
# ./scripts/domestic_full.sh --resume # 从中断处继续(跳过已完成步骤,
|
||||
# # 步骤状态见 data/run_state/{date}.state)
|
||||
# =============================================
|
||||
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 "══════ 国内全流程开始 ══════"
|
||||
LOG "══════ 国内全流程开始(resume=${RESUME})═══════"
|
||||
|
||||
# ── 加载 .env ──
|
||||
export $(grep -v '^#' .env | grep -v '^$' | xargs 2>/dev/null || true)
|
||||
|
||||
# ── 1. M1 Pi 抓取(headful Playwright + HTTP 代理)──
|
||||
LOG "[1/2] Pi M1 抓取(8G headful + HTTP 代理)..."
|
||||
bash "$SCRIPT_DIR/domestic_crawl_8g.sh" 2>&1 | tail -5 || LOG "WARNING: 部分源抓取失败,继续管道"
|
||||
# 部分源抓取失败不阻塞管道(原语义);抓取本身按 URL 去重(index.jsonl),
|
||||
# 已抓取过的 URL 不会重复写入
|
||||
if step_should_run M1_crawl; then
|
||||
LOG "[1/2] Pi M1 抓取(8G headful + HTTP 代理)..."
|
||||
bash "$SCRIPT_DIR/domestic_crawl_8g.sh" 2>&1 | tail -5 || LOG "WARNING: 部分源抓取失败,继续管道"
|
||||
step_mark M1_crawl
|
||||
fi
|
||||
|
||||
# ── 2. M2→M6 管道(含日报)──
|
||||
LOG "[2/2] 全链路管道..."
|
||||
bash "$SCRIPT_DIR/pipeline.sh" 2>&1 | tail -10
|
||||
if [ "$RESUME" = "1" ]; then
|
||||
bash "$SCRIPT_DIR/pipeline.sh" --resume 2>&1 | tail -10
|
||||
else
|
||||
bash "$SCRIPT_DIR/pipeline.sh" 2>&1 | tail -10
|
||||
fi
|
||||
|
||||
LOG "══════ 国内全流程完成 ✅ ══════"
|
||||
|
||||
+32
-15
@@ -2,9 +2,19 @@
|
||||
# =============================================
|
||||
# 国内服务器:全链路管道 M2 → M3 → M4 → M5 → M6 → 日报
|
||||
# =============================================
|
||||
# 用法:./scripts/pipeline.sh
|
||||
# 用法:
|
||||
# ./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)"
|
||||
@@ -12,57 +22,64 @@ PROJECT_DIR="$(dirname "$SCRIPT_DIR")"
|
||||
|
||||
cd "$PROJECT_DIR"
|
||||
|
||||
echo "[$(date)] ═══ 全链路管道开始 ═══"
|
||||
# ── 参数解析 ──
|
||||
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 "══════ 全链路管道开始(resume=${RESUME})═══════"
|
||||
|
||||
# 加载 .env
|
||||
export $(grep -v '^#' .env | grep -v '^$' | xargs 2>/dev/null || true)
|
||||
|
||||
# ── M2: 正文提取 ──
|
||||
echo "[$(date)] [M2] 正文提取..."
|
||||
.venv/bin/python3 -c "
|
||||
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: 去重 ──
|
||||
echo "[$(date)] [M3] 三层去重..."
|
||||
.venv/bin/python3 -c "
|
||||
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: 翻译+事件 ──
|
||||
echo "[$(date)] [M4] 翻译+事件抽取..."
|
||||
.venv/bin/python3 -c "
|
||||
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: 向量生成 ──
|
||||
echo "[$(date)] [M5] 向量生成..."
|
||||
.venv/bin/python3 -c "
|
||||
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 入库 ──
|
||||
echo "[$(date)] [M6] Qdrant 入库..."
|
||||
.venv/bin/python3 -c "
|
||||
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')
|
||||
"
|
||||
|
||||
# ── 日报 ──
|
||||
echo "[$(date)] [日报] 生成日报(结构化入库)..."
|
||||
.venv/bin/python3 -c "
|
||||
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 '日报: 无数据/失败')
|
||||
"
|
||||
|
||||
echo "[$(date)] ═══ 全链路管道完成 ✅ ═══"
|
||||
LOG "══════ 全链路管道完成 ✅ ═══════"
|
||||
|
||||
Reference in New Issue
Block a user