From 06b00b7f492e8afde1b4a91108ce5957c5c3eeac Mon Sep 17 00:00:00 2001 From: Simon Date: Wed, 12 Aug 2026 11:06:29 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=85=A8=E6=B5=81=E7=A8=8B=E8=84=9A?= =?UTF-8?q?=E6=9C=AC=20--resume=20=E6=96=AD=E7=82=B9=E7=BB=AD=E8=B7=91=20+?= =?UTF-8?q?=20=E5=90=84=E6=AD=A5=E9=AA=A4=E6=96=87=E4=BB=B6=E7=BA=A7?= =?UTF-8?q?=E5=A2=9E=E9=87=8F=E8=AF=B4=E6=98=8E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增 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 机制正常) --- .gitignore | 1 + README.md | 14 ++++++++++ continuation.md | 21 ++++++++++++++ scripts/_step_state.sh | 60 ++++++++++++++++++++++++++++++++++++++++ scripts/domestic_full.sh | 33 +++++++++++++++++++--- scripts/pipeline.sh | 47 +++++++++++++++++++++---------- 6 files changed, 157 insertions(+), 19 deletions(-) create mode 100644 scripts/_step_state.sh diff --git a/.gitignore b/.gitignore index dff8073..ee18f35 100644 --- a/.gitignore +++ b/.gitignore @@ -31,6 +31,7 @@ data/events/ data/embeddings/ data/qdrant_storage/ data/reports/ +data/run_state/ # Logs logs/*.log diff --git a/README.md b/README.md index 7beb861..3776888 100644 --- a/README.md +++ b/README.md @@ -104,6 +104,12 @@ sudo systemctl enable --now privoxy # 全流程:M1 抓取 → M2→M6 管道 → 日报 # 每天 06:00 / 12:00 / 18:00 / 22:00 自动执行 bash scripts/domestic_full.sh + +# 中断恢复:跳过当天已完成步骤(状态存 data/run_state/{date}.state) +# 支持步骤级断点:M1_crawl / M2_extract / M3_dedup / M4_translate / M5_embed / M6_index / report +bash scripts/domestic_full.sh --resume +# 仅 M2→M6 管道(同样支持 --resume) +bash scripts/pipeline.sh --resume ``` ### 单步执行 @@ -161,6 +167,14 @@ ls -lt data/reports/ 去重时跨源重复的新闻,会把所有来源记录到保留的唯一篇 `source_ids` 字段(首个来源为 `source_id`),经翻译透传后在日报事件 `source` 展示(如 "Barron's, CNBC, Reuters",最多 3 个)。 +### 中断恢复(--resume) + +全流程脚本支持断点续跑: + +- `scripts/domestic_full.sh --resume` / `scripts/pipeline.sh --resume`:跳过当天已完成步骤(步骤状态存 `data/run_state/{YYYYMMDD}.state`,按天换新,跨天自动失效) +- 步骤粒度:`M1_crawl` / `M2_extract` / `M3_dedup` / `M4_translate` / `M5_embed` / `M6_index` / `report`;某步失败不标记,`--resume` 从失败处重试 +- **文件级增量(各步骤内自动跳过已处理文件)**:M2 按 `data/processed/.../{url_hash}.json` 存在性、M3 按指纹库判重、M4 按 `data/events/.../{url_hash}.json` 存在性、M5 按 `data/embeddings/.../index.json`、M6 按 Qdrant upsert 幂等(点 id=url_hash)、日报按 MySQL 唯一键覆盖 + ### 日报入库(M9) 日报内容结构化写入与 [news 项目](https://github.com/) 共用的 MySQL `myquant` 库(表 `news_report` / `news_event`,`report_type="intl"`,同一天重复生成幂等覆盖)。表结构与数据契约见 news 项目 `docs/db_schema.md`。 diff --git a/continuation.md b/continuation.md index badd66e..1861d40 100644 --- a/continuation.md +++ b/continuation.md @@ -4,6 +4,27 @@ --- +## 2026-08-12 会话成果(二):全流程 --resume 断点续跑 + +**背景:** `scripts/domestic_full.sh`(M1→M6→日报)中断(ssh 断线/断电)后只能从头重跑。 + +**改动:** + +| 文件 | 改动内容 | +|------|---------| +| `scripts/_step_state.sh`(新增) | 步骤状态库:`step_run` / `step_should_run` / `step_mark` / `step_done`;状态文件 `data/run_state/{YYYYMMDD}.state`(按天换新);`step_run` 显式检查退出码(规避 bash set -e 条件上下文陷阱) | +| `scripts/pipeline.sh` | 支持 `--resume`;M2→M6→report 每步 `step_run` 包裹 | +| `scripts/domestic_full.sh` | 支持 `--resume`;M1 用 `step_should_run`/`step_mark`(部分源失败仍标记,原语义);透传 `--resume` 给 pipeline.sh | +| `.gitignore` | 忽略 `data/run_state/` | + +**文件级增量确认(各步骤已内置,无需新写):** M2 processed JSON 存在跳过、M3 指纹库判重、M4 events JSON 存在跳过、M5 embedding index、M6 Qdrant upsert 幂等(id=url_hash)、日报 MySQL 唯一键覆盖。 + +**测试:** 本地模拟 3 场景全过(失败步骤不 mark / set -e 退出不 mark / resume 跳过已完成);`--badarg` exit=1;pi5 同步后验证 step_run 机制与参数校验正常。 + +**注意:** crontab 07/12/18 跑的是不带 --resume 的全量(步骤级不跳过,靠文件级增量),--resume 供手动中断恢复。 + +--- + ## 2026-08-12 会话成果 ### M9.2:AI 模型按场景配置 + 去重多来源 diff --git a/scripts/_step_state.sh b/scripts/_step_state.sh new file mode 100644 index 0000000..ed97f35 --- /dev/null +++ b/scripts/_step_state.sh @@ -0,0 +1,60 @@ +#!/bin/bash +# ============================================= +# 步骤状态管理(--resume 中断恢复用) +# ============================================= +# 用法(在脚本中 source): +# source "$SCRIPT_DIR/_step_state.sh" +# RESUME=0|1 # 由入口脚本参数解析后设置 +# step_run <步骤名> # 执行并在成功后标记;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" +} diff --git a/scripts/domestic_full.sh b/scripts/domestic_full.sh index ae4b06f..9e405ea 100755 --- a/scripts/domestic_full.sh +++ b/scripts/domestic_full.sh @@ -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 "══════ 国内全流程完成 ✅ ══════" diff --git a/scripts/pipeline.sh b/scripts/pipeline.sh index ffbd2a0..832f4b1 100755 --- a/scripts/pipeline.sh +++ b/scripts/pipeline.sh @@ -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 "══════ 全链路管道完成 ✅ ═══════"