diff --git a/continuation.md b/continuation.md index 5a4e1da..266cb05 100644 --- a/continuation.md +++ b/continuation.md @@ -4,6 +4,29 @@ --- +## 本次完成 (2026-08-22) — P1-1 修复 dedup 高重复率返回码语义 + +**问题**:`run_dedup` 重复率 > 5% 时返回 1;`scheduler/pipeline.py` 曾把 `dedup` 的 rc=1 无条件视为成功并打日志"无新数据场景"。结果:高重复率(可能是正常无新数据,也可能是抓取源/指纹库异常)被一刀切掩盖,真实异常无法上报。 + +**修复(分离"执行成功"与"统计告警")**: +- `scripts/run_dedup.py`: + - 生产模式(默认):重复率仅作 **WARNING 告警**,不影响退出码(执行成功即 0);阈值常量 `_DUP_RATE_THRESHOLD = 0.05` + - 新增 `--strict` 验收模式:保留 M3 验收门槛,重复率 > 5% 时返回 1(供人工验收) + - 新增统计快照 `data/deduped/{date}/stats.json`(原子写):unique/duplicates/total/dup_rate/layers/fingerprint_total/generated_at,供运维排查与监控 + - 顺带修复:空日场景 `sources.json` 写入前补 `mkdir`(原会 FileNotFoundError) +- `scheduler/pipeline.py`:移除 `dedup` rc=1 特判,恢复"非 0 即失败"的统一语义 + +**测试**: +- 新增 `tests/test_run_dedup.py` 7 个(生产模式高重复返回 0 / --strict 高重复返回 1 / --strict 低重复返回 0 / 空日返回 0 / stats.json 结构与数值 / pipeline 不再掩蔽失败) +- 全量 **266 passed**,3 个 crawler 基线失败与本次无关 + +**生产实测**:同日重跑重复率 100% → 生产模式 `rc=0` + WARNING + stats.json 完整;`--strict` 同场景 `rc=1`(用非管道方式核实退出码) + +**待办/遗留**: +- P1-2(守护进程补跑 steps 含 cninfo 三步)、P1-3(调度时区与 date.today() 一致性),待用户决策 + +--- + ## 本次完成 (2026-08-22) — P0-3/P0-4 补跑日期语义修复 **P0-3(crawler 补跑历史日期静默空转)**: @@ -20,7 +43,7 @@ - 全量 **259 passed**,3 个 crawler 基线失败与本次无关 **待办/遗留**: -- P1:dedup rc=1 静默转成功(天天重复率>5% 被掩蔽)、守护进程补跑 steps 含 cninfo_*、调度时区与 date.today() 一致性,待用户决策 +- ~~P1-1(dedup rc=1 静默转成功)~~ 已修复(见上方小节);余:守护进程补跑 steps 含 cninfo_*、调度时区与 date.today() 一致性,待用户决策 - 既有小问题:`tests/test_scheduler.py` 部分 run_pipeline 测试未传 state_path,会写真实 `data/pipeline/state.json`(既有行为,未改) --- diff --git a/scheduler/pipeline.py b/scheduler/pipeline.py index ec8c7e0..a025162 100644 --- a/scheduler/pipeline.py +++ b/scheduler/pipeline.py @@ -290,11 +290,6 @@ def run_step(name: str, date_str: str) -> StepResult: elapsed = (datetime.now() - started).total_seconds() ok = proc.returncode == 0 - # dedup 返回 1 是"重复率过高"(无新文章的正常场景) - if name == "dedup" and proc.returncode == 1: - ok = True - logger.info("dedup 重复率超过阈值(无新数据场景,视为成功)") - tail_msg = f"rc={proc.returncode}" if not ok else "" if ok: diff --git a/scripts/run_dedup.py b/scripts/run_dedup.py index c84e93d..ec9b9d2 100644 --- a/scripts/run_dedup.py +++ b/scripts/run_dedup.py @@ -22,7 +22,7 @@ import argparse import json import sys from collections import Counter -from datetime import date +from datetime import date, datetime from pathlib import Path from loguru import logger @@ -31,6 +31,9 @@ from pydantic import ValidationError from dedup import Deduper from extractor import Article +# 重复率告警阈值(超过则输出 WARNING 日志;--strict 验收模式下作为返回码门槛) +_DUP_RATE_THRESHOLD = 0.05 + def _setup_logger(level: str) -> None: logger.remove() @@ -200,6 +203,11 @@ def main() -> int: parser.add_argument("--simhash-threshold", type=int, default=3) parser.add_argument("--window-days", type=int, default=30) parser.add_argument("--reset", action="store_true", help="处理前清空指纹库") + parser.add_argument( + "--strict", action="store_true", + help="验收模式:重复率超过阈值(5%)时返回 1。" + "默认为生产模式,重复率仅作统计告警,不影响退出码", + ) parser.add_argument("--log-level", default="INFO") args = parser.parse_args() @@ -243,6 +251,7 @@ def main() -> int: # {url_hash: [source_id, ...]},配合 uniques/{url_hash}.json 的 sources 字段 # 与指纹库 source_ids 列,提供「一条唯一新闻多个来源」的完整记录。 sources_path = out_root / args.date / "sources.json" + sources_path.parent.mkdir(parents=True, exist_ok=True) sources_path.write_text( json.dumps( {k: v for k, v in sources_map.items() if v}, @@ -260,15 +269,47 @@ def main() -> int: total = total_uniq + total_dup rate = total_dup / max(total, 1) - logger.info( + log_fn = logger.warning if rate > _DUP_RATE_THRESHOLD else logger.info + log_fn( "全部完成: 唯一 {} / 重复 {} (重复率 {:.1%}) layers={}", total_uniq, total_dup, rate, dict(total_layers), ) - # 验收门槛: ≤ 5% - return 0 if rate <= 0.05 or total == 0 else 1 + + # 统计快照: data/deduped/{date}/stats.json(原子写),供运维排查与监控 + stats_payload = { + "date": args.date, + "generated_at": datetime.now().isoformat(), + "unique": total_uniq, + "duplicates": total_dup, + "total": total, + "dup_rate": round(rate, 4), + "dup_rate_threshold": _DUP_RATE_THRESHOLD, + "layers": dict(total_layers), + "fingerprint_total": deduper.stats().total, + } + stats_path = out_root / args.date / "stats.json" + stats_path.parent.mkdir(parents=True, exist_ok=True) + tmp = stats_path.with_suffix(".json.tmp") + tmp.write_text(json.dumps(stats_payload, ensure_ascii=False, indent=2), + encoding="utf-8") + tmp.replace(stats_path) + + if rate > _DUP_RATE_THRESHOLD and total > 0: + logger.warning( + "重复率 {:.1%} 超过阈值 {:.0%}(唯一 {},重复 {})。" + "同日多次调度/跨源转载属常见现象;若为当日首次处理仍异常偏高," + "请检查指纹库与抓取源(详见 {})", + rate, _DUP_RATE_THRESHOLD, total_uniq, total_dup, stats_path, + ) + + # 生产模式:重复率是统计指标,不影响退出码(执行成功即 0); + # 验收模式(--strict):保留 M3 验收门槛(重复率 ≤ 5%),超标返回 1。 + if args.strict and total > 0 and rate > _DUP_RATE_THRESHOLD: + return 1 + return 0 if __name__ == "__main__": diff --git a/tests/test_run_dedup.py b/tests/test_run_dedup.py new file mode 100644 index 0000000..24ecc47 --- /dev/null +++ b/tests/test_run_dedup.py @@ -0,0 +1,173 @@ +"""run_dedup 脚本测试(生产/验收返回码语义 + stats.json 统计快照)。 + +覆盖: + - 生产模式(默认):高重复率不影响退出码,返回 0;重复率仅作告警 + - 验收模式(--strict):高重复率返回 1(保留 M3 验收门槛 ≤ 5%) + - 低重复率场景两种模式均返回 0 + - stats.json 统计快照结构正确(原子写),含重复率与指纹库总量 +""" + +from __future__ import annotations + +import hashlib +import json +import sys +from pathlib import Path + +import pytest + + +def _url_hash(url: str) -> str: + return hashlib.sha1(url.encode("utf-8")).hexdigest()[:16] + + +def _write_article(proc_dir: Path, url: str, title: str, content: str) -> None: + """写一个 Article JSON,文件名为 url_hash。""" + h = _url_hash(url) + article = { + "source_id": "testsrc", + "url": url, + "url_hash": h, + "title": title, + "content": content, + "word_count": len(content), + } + (proc_dir / f"{h}.json").write_text( + json.dumps(article, ensure_ascii=False), encoding="utf-8" + ) + + +def _setup_fixture(tmp_path: Path, *, dup_count: int, uniq_count: int) -> Path: + """构造 processed 夹具:uniq_count 篇互不相同 + dup_count 篇与第 1 篇内容相同。""" + day = "20260616" + proc_dir = tmp_path / "processed" / "testsrc" / day + proc_dir.mkdir(parents=True) + + base_content = "这是一条足够长的测试新闻正文内容,用于触发去重判定逻辑。" * 3 + # 有重复篇时才写基准篇(重复目标);否则唯一数 = uniq_count + if dup_count > 0: + _write_article(proc_dir, "https://x.example/base", "基准新闻", base_content) + # 与基准篇内容相同、URL 不同 → 判为重复 + for i in range(dup_count): + _write_article( + proc_dir, f"https://x.example/dup{i}", f"重复新闻{i}", base_content + ) + # 互不相同的其它唯一文章 + for i in range(uniq_count): + _write_article( + proc_dir, f"https://x.example/uniq{i}", f"独立新闻{i}", + f"完全不同的正文内容片段编号 {i},讲述另一件事。" * 3, + ) + return proc_dir + + +def _run_dedup_main(tmp_path: Path, *extra_args: str) -> int: + """以指定参数调用 scripts.run_dedup.main,返回退出码。""" + import scripts.run_dedup as mod + + day = "20260616" + argv = [ + "run_dedup", + "--processed-root", str(tmp_path / "processed"), + "--out-root", str(tmp_path / "deduped"), + "--db", str(tmp_path / "fingerprints.sqlite3"), + "--date", day, + "--log-level", "ERROR", + *extra_args, + ] + old_argv = sys.argv + sys.argv = argv + try: + return mod.main() + finally: + sys.argv = old_argv + + +def _read_stats(tmp_path: Path) -> dict: + p = tmp_path / "deduped" / "20260616" / "stats.json" + assert p.is_file(), "stats.json 未生成" + return json.loads(p.read_text(encoding="utf-8")) + + +# --------------------------------------------------------------------------- # +# 返回码语义 +# --------------------------------------------------------------------------- # + +def test_production_mode_high_dup_returns_0(tmp_path: Path) -> None: + """生产模式:高重复率(90%)不影响退出码,返回 0。""" + _setup_fixture(tmp_path, dup_count=9, uniq_count=0) + rc = _run_dedup_main(tmp_path) + assert rc == 0, "生产模式下重复率仅作告警,不应改变退出码" + + +def test_strict_mode_high_dup_returns_1(tmp_path: Path) -> None: + """验收模式 --strict:重复率超过 5% 门槛时返回 1。""" + _setup_fixture(tmp_path, dup_count=9, uniq_count=0) + rc = _run_dedup_main(tmp_path, "--strict") + assert rc == 1 + + +def test_strict_mode_low_dup_returns_0(tmp_path: Path) -> None: + """验收模式:重复率 ≤ 5% 时返回 0。""" + _setup_fixture(tmp_path, dup_count=0, uniq_count=10) + rc = _run_dedup_main(tmp_path, "--strict") + assert rc == 0 + + +def test_empty_day_returns_0(tmp_path: Path) -> None: + """当日无文章时返回 0(两种模式)。""" + proc_dir = tmp_path / "processed" / "testsrc" / "20260616" + proc_dir.mkdir(parents=True) + assert _run_dedup_main(tmp_path) == 0 + assert _run_dedup_main(tmp_path, "--strict") == 0 + + +# --------------------------------------------------------------------------- # +# stats.json 统计快照 +# --------------------------------------------------------------------------- # + +def test_stats_json_written_and_correct(tmp_path: Path) -> None: + """stats.json 结构正确:重复率、唯一/重复数、指纹库总量。""" + _setup_fixture(tmp_path, dup_count=9, uniq_count=0) + _run_dedup_main(tmp_path) + s = _read_stats(tmp_path) + assert s["date"] == "20260616" + assert s["unique"] == 1 + assert s["duplicates"] == 9 + assert s["total"] == 10 + assert s["dup_rate"] == pytest.approx(0.9, abs=1e-6) + assert s["dup_rate_threshold"] == pytest.approx(0.05) + # 指纹库只保留唯一文章 + assert s["fingerprint_total"] == 1 + assert "generated_at" in s and "layers" in s + + +def test_stats_json_low_dup_rate(tmp_path: Path) -> None: + """低重复率场景:dup_rate 为 0。""" + _setup_fixture(tmp_path, dup_count=0, uniq_count=5) + _run_dedup_main(tmp_path) + s = _read_stats(tmp_path) + assert s["unique"] == 5 + assert s["duplicates"] == 0 + assert s["dup_rate"] == 0.0 + assert s["fingerprint_total"] == 5 + + +# --------------------------------------------------------------------------- # +# pipeline: dedup 不再特判(真实失败现在会正确上报) +# --------------------------------------------------------------------------- # + +def test_pipeline_dedup_failure_no_longer_masked(monkeypatch: pytest.MonkeyPatch) -> None: + """移除特判后:dedup 非 0 退出码应如实标记为失败。""" + from datetime import date + from types import SimpleNamespace + + from scheduler import pipeline + + def _fake_run(cmd, timeout=None): # noqa: ARG001 + return SimpleNamespace(returncode=1) + + monkeypatch.setattr(pipeline.subprocess, "run", _fake_run) + sr = pipeline.run_step("dedup", date.today().strftime("%Y%m%d")) + assert sr.success is False, "dedup 返回非 0 应标记失败(特判已移除)" + assert "rc=1" in sr.tail_msg