feat: 增量处理与 pipeline 断点续跑

- M2 run_extractor: 产物存在即跳过提取(仅回补 index),--force 全量;
  增量成功率统计含跳过项,修复全跳过时误报失败
- M4 run_event_extraction: data/events/{day}/{url_hash}.json 已存在即跳过,
  不重复调用 LLM API;--force 全量;failed 只留本次失败
- M5 run_embedding: data/embeddings/{day}/{url_hash}.json 已存在即跳过,
  不重复调用 embed API;--force 全量
- scheduler/pipeline: 步骤结果按日期写入 data/pipeline/state.json(原子写),
  run_pipeline(resume=True) 从首个失败/未执行步骤续跑
- run_scheduler + a-share CLI: --once --resume 断点续跑(--steps 互斥)
- 新增 tests/test_incremental.py 10 个测试(跳过逻辑 + resume)
- .gitignore: 忽略 data/pipeline/ 运行状态
This commit is contained in:
2026-08-12 10:19:27 +08:00
parent c1a803968a
commit 2b4efea219
8 changed files with 518 additions and 35 deletions
+49 -8
View File
@@ -13,6 +13,9 @@
data/embeddings/{day}/{url_hash}.json (含 vector 完整内容)
data/embeddings/{day}/index.jsonl (扁平摘要,不含向量,便于检索/调试)
data/embeddings/{day}/failed.jsonl (失败列表)
增量: 默认跳过已嵌入的文章(输出目录已有 {url_hash}.json 视为已处理),
断点续跑/失败重试不会重复调用 embed API;--force 强制全量重嵌入。
"""
from __future__ import annotations
@@ -142,6 +145,26 @@ def _collect_inputs(args: argparse.Namespace) -> list[tuple[Path, str]]:
return files
def _filter_existing(
files: list[tuple[Path, str]], out_dir: Path
) -> tuple[list[tuple[Path, str]], int]:
"""过滤掉已有产物(输出目录存在同名 {url_hash}.json)的输入。
输入与输出文件名均为 {url_hash}.json,直接比对 stem。
返回 (待处理, 跳过数);断点续跑/失败重试借此避免重复调用 embed API。
"""
pending: list[tuple[Path, str]] = []
skipped = 0
for fp, kind in files:
if (out_dir / f"{fp.stem}.json").exists():
skipped += 1
else:
pending.append((fp, kind))
if skipped:
logger.info("跳过已嵌入 {} 篇(产物已存在),待处理 {}", skipped, len(pending))
return pending, skipped
# --------------------------------------------------------------------------- #
# 主流程
# --------------------------------------------------------------------------- #
@@ -150,12 +173,24 @@ async def _run(args: argparse.Namespace) -> int:
load_dotenv()
files = _collect_inputs(args)
if args.limit:
files = files[: args.limit]
if not files:
logger.error("未发现任何输入文件: {} ({})", args.input, args.date)
return 2
logger.info("待嵌入文章数: {} (input={})", len(files), args.input)
out_dir = Path(args.out_root) / args.date
out_dir.mkdir(parents=True, exist_ok=True)
# 增量:跳过已有产物(断点续跑/失败重试不重复调用 embed API),--force 全量
skipped = 0
if not args.force:
files, skipped = _filter_existing(files, out_dir)
if args.limit:
files = files[: args.limit]
if not files:
logger.info(
"无待嵌入文章(全部已处理,跳过 {} 篇),如需重新嵌入请加 --force", skipped
)
return 0
logger.info("待嵌入文章数: {} (跳过已处理 {}; input={})", len(files), skipped, args.input)
# 准备每篇文本
prepared: list[tuple[str, Article, str | None]] = []
@@ -170,13 +205,17 @@ async def _run(args: argparse.Namespace) -> int:
logger.error("所有输入文件均无法解析")
return 2
out_dir = Path(args.out_root) / args.date
out_dir.mkdir(parents=True, exist_ok=True)
index_path = out_dir / "index.jsonl"
failed_path = out_dir / "failed.jsonl"
for p in (index_path, failed_path):
if p.exists():
p.unlink()
if args.force:
# 全量模式:重建 index / failed
for p in (index_path, failed_path):
if p.exists():
p.unlink()
else:
# 增量模式:index 累积追加;failed 只保留本次运行失败的
if failed_path.exists():
failed_path.unlink()
started = time.time()
succ_cnt = 0
@@ -282,6 +321,8 @@ def main() -> int:
help="每批送 embed 的条数(DashScope 上限 10)")
parser.add_argument("--limit", type=int, default=0,
help="最多处理 N 篇,0=不限")
parser.add_argument("--force", action="store_true",
help="强制全量重嵌入(默认跳过已嵌入文章)")
parser.add_argument("--log-level", default="INFO")
args = parser.parse_args()
+47 -7
View File
@@ -6,12 +6,16 @@
data/events/{YYYYMMDD}/index.jsonl (扁平摘要)
data/events/{YYYYMMDD}/failed.jsonl (失败列表)
增量: 默认跳过已抽取的文章(输出目录已有 {url_hash}.json 视为已处理),
断点续跑/失败重试不会重复调用 LLM API;--force 强制全量重抽。
用法:
uv run python -m scripts.run_event_extraction
uv run python -m scripts.run_event_extraction --date 20260616
uv run python -m scripts.run_event_extraction --provider qwen --model qwen-plus
uv run python -m scripts.run_event_extraction --concurrency 5 --limit 10
uv run python -m scripts.run_event_extraction --input-root data/processed --no-deduped
uv run python -m scripts.run_event_extraction --force # 全量重抽
"""
from __future__ import annotations
@@ -79,6 +83,24 @@ def _collect_inputs(
return files
def _filter_existing(files: list[Path], out_dir: Path) -> tuple[list[Path], int]:
"""过滤掉已有产物(输出目录存在同名 {url_hash}.json)的输入。
输入文件名即 url_hash(如 {url_hash}.json),与 M4 产物命名一致。
返回 (待处理文件, 跳过数);断点续跑/失败重试借此避免重复调用 LLM API。
"""
pending: list[Path] = []
skipped = 0
for fp in files:
if (out_dir / f"{fp.stem}.json").exists():
skipped += 1
else:
pending.append(fp)
if skipped:
logger.info("跳过已处理 {} 篇(产物已存在),待处理 {}", skipped, len(pending))
return pending, skipped
def _load_article(p: Path) -> Article | None:
try:
return Article.model_validate(json.loads(p.read_text(encoding="utf-8")))
@@ -101,24 +123,40 @@ async def _run(args: argparse.Namespace) -> int:
input_root = Path(args.input_root)
use_deduped = not args.no_deduped
files = _collect_inputs(input_root, args.date, use_deduped, args.source)
if args.limit:
files = files[: args.limit]
if not files:
logger.error(
"{} 下未发现 {} 的文章(use_deduped={})",
input_root, args.date, use_deduped,
)
return 2
logger.info("待处理文章数: {}", len(files))
out_dir = Path(args.out_root) / args.date
out_dir.mkdir(parents=True, exist_ok=True)
# 增量:跳过已有产物(断点续跑/失败重试不重复调用 LLM API),--force 全量
skipped = 0
if not args.force:
files, skipped = _filter_existing(files, out_dir)
if args.limit:
files = files[: args.limit]
if not files:
logger.info(
"无待处理文章(全部已抽取,跳过 {} 篇),如需重抽请加 --force", skipped
)
return 0
logger.info("待处理文章数: {} (跳过已处理 {})", len(files), skipped)
index_path = out_dir / "index.jsonl"
failed_path = out_dir / "failed.jsonl"
# 重跑时清掉旧的 jsonl,避免重复追加
for p in (index_path, failed_path):
if p.exists():
p.unlink()
if args.force:
# 全量模式:重建 index / failed
for p in (index_path, failed_path):
if p.exists():
p.unlink()
else:
# 增量模式:index 累积追加;failed 只保留本次运行失败的
if failed_path.exists():
failed_path.unlink()
template = PromptTemplate(args.prompt)
semaphore = asyncio.Semaphore(args.concurrency)
@@ -209,6 +247,8 @@ def main() -> int:
help="LLM 异步并发上限")
parser.add_argument("--max-attempts", type=int, default=3,
help="单篇文章最大重试次数")
parser.add_argument("--force", action="store_true",
help="强制全量重抽(默认跳过已抽取文章)")
parser.add_argument("--limit", type=int, default=0,
help="最多处理 N 篇,0=不限制(用于联调)")
parser.add_argument("--prompt", default="prompts/event_extraction.md",
+51 -16
View File
@@ -14,12 +14,14 @@
from __future__ import annotations
import argparse
import contextlib
import json
import sys
from datetime import date, datetime
from pathlib import Path
from loguru import logger
from pydantic import ValidationError
from extractor import Article, ExtractError, extract_article
from extractor.parser import _url_hash
@@ -238,16 +240,21 @@ def _process_cninfo_v2(rec: dict, json_path: Path, out_dir: Path) -> Article | N
return _save_article(article, out_dir)
def _append_index(article: Article, out_dir: Path) -> None:
"""把 Article 的扁平摘要追加到 index.jsonl。"""
flat = article.model_dump(exclude={"content", "images"}, mode="json")
flat["article_file"] = f"{article.url_hash}.json"
flat["content_preview"] = article.content[:80]
with (out_dir / "index.jsonl").open("a", encoding="utf-8") as f:
f.write(json.dumps(flat, ensure_ascii=False) + "\n")
def _save_article(article: Article, out_dir: Path) -> Article:
"""保存 Article JSON 并追加 index。"""
out_dir.mkdir(parents=True, exist_ok=True)
article_path = out_dir / f"{article.url_hash}.json"
article_path.write_text(article.model_dump_json(indent=2), encoding="utf-8")
flat = article.model_dump(exclude={"content", "images"}, mode="json")
flat["article_file"] = article_path.name
flat["content_preview"] = article.content[:80]
with (out_dir / "index.jsonl").open("a", encoding="utf-8") as f:
f.write(json.dumps(flat, ensure_ascii=False) + "\n")
_append_index(article, out_dir)
return article
@@ -257,37 +264,57 @@ def _process_source_day(
raw_root: Path,
out_root: Path,
body_xpath_map: dict[str, str] | None = None,
) -> tuple[int, int]:
"""处理单个源单日。返回 (成功数, 总数)。"""
*,
force: bool = False,
) -> tuple[int, int, int]:
"""处理单个源单日。返回 (成功数, 总数, 跳过数)。
默认增量:已提取的文章(输出目录已有 {url_hash}.json)跳过提取,仅回补 index 行;
force=True 时全量重提取并重建 index。
"""
raw_dir = raw_root / source_id / day
out_dir = out_root / source_id / day
records = _iter_article_records(raw_dir)
if not records:
logger.info("源 {} 日期 {} 无可处理记录", source_id, day)
return 0, 0
return 0, 0, 0
# 清理同日旧的 index.jsonl,避免重复追加
# 全量模式:重建 index;增量模式:保留旧 index 追加新条目
old_index = out_dir / "index.jsonl"
if old_index.exists():
if force and old_index.exists():
old_index.unlink()
succ = 0
skipped = 0
for rec in records:
url_hash = rec.get("url_hash") or _url_hash(rec.get("url") or "")
existing = out_dir / f"{url_hash}.json"
if not force and existing.is_file():
# 增量:跳过已提取,回补 index 行保持摘要完整
skipped += 1
with contextlib.suppress(json.JSONDecodeError, ValidationError, OSError):
_append_index(
Article.model_validate(json.loads(existing.read_text(encoding="utf-8"))),
out_dir,
)
continue
article = _process_one(rec, raw_dir, out_dir, body_xpath_map)
if article is not None:
succ += 1
total = len(records)
rate = succ / max(total, 1)
# 增量模式下「跳过已提取」视为已成功处理,避免全跳过时误报成功率 0%
rate = (succ + skipped) / max(total, 1)
logger.info(
"源 {} 日期 {} 提取完成: {}/{} 成功率 {:.0%}",
"源 {} 日期 {} 提取完成: {}/{} 成功率 {:.0%} (跳过已提取 {})",
source_id,
day,
succ,
total,
rate,
skipped,
)
return succ, total
return succ, total, skipped
def _list_source_dirs(raw_root: Path) -> list[str]:
@@ -307,6 +334,8 @@ def main() -> int:
default=date.today().strftime("%Y%m%d"),
help="处理日期 YYYYMMDD,默认今日",
)
parser.add_argument("--force", action="store_true",
help="强制全量重提取(默认跳过已提取文章)")
parser.add_argument("--log-level", default="INFO")
args = parser.parse_args()
@@ -335,18 +364,24 @@ def main() -> int:
started = datetime.now()
total_succ = 0
total_all = 0
total_skipped = 0
for src in sources:
succ, total = _process_source_day(src, args.date, raw_root, out_root, body_xpath_map)
succ, total, skipped = _process_source_day(
src, args.date, raw_root, out_root, body_xpath_map, force=args.force
)
total_succ += succ
total_all += total
total_skipped += skipped
elapsed = (datetime.now() - started).total_seconds()
rate = total_succ / max(total_all, 1)
# 增量模式下跳过已提取视为成功
rate = (total_succ + total_skipped) / max(total_all, 1)
logger.info(
"全部完成: {}/{} 成功率 {:.0%} 用时 {:.1f}s",
"全部完成: {}/{} 成功率 {:.0%} (跳过已提取 {}) 用时 {:.1f}s",
total_succ,
total_all,
rate,
total_skipped,
elapsed,
)
return 0 if rate >= 0.9 or total_all == 0 else 1
+12 -2
View File
@@ -59,11 +59,17 @@ def _parse_schedule_times(raw: str) -> list[tuple[int, int]]:
def _once(args: argparse.Namespace) -> int:
"""单次执行模式。"""
"""单次执行模式。
默认全量执行;--resume 时断点续跑(跳过连续成功步骤,从失败/未执行步骤继续)。
"""
steps = None
if args.steps:
steps = [s.strip() for s in args.steps.split(",")]
run_pipeline(args.date, steps=steps)
if args.resume and args.steps:
logger.error("--resume 与 --steps 不能同时使用(断点续跑针对全链路)")
return 2
run_pipeline(args.date, steps=steps, resume=args.resume)
return 0
@@ -185,6 +191,10 @@ def main() -> int:
)
parser.add_argument("--steps", default=None,
help="仅执行指定步骤,逗号分隔 (如 crawler,extractor)")
parser.add_argument(
"--resume", action="store_true",
help="断点续跑(仅 --once):跳过连续成功步骤,从上次失败/未执行步骤继续",
)
parser.add_argument("--log-level", default="INFO")
args = parser.parse_args()