From 7ea8925209a958c6a20ae91fb1b4377086f66329 Mon Sep 17 00:00:00 2001 From: simon Date: Sat, 22 Aug 2026 19:11:32 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=20xwlb=20=E6=97=A5?= =?UTF-8?q?=E6=9C=9F=E9=94=99=E4=BD=8D=E4=B8=8E=E9=87=8D=E5=A4=8D=E8=A1=8C?= =?UTF-8?q?=E9=97=AE=E9=A2=98(=E6=96=B9=E6=A1=88=20A)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - run_xwlb: 新增 --date(处理日)参数,落盘改为处理日目录,下游 M2-M6 打通 - run_xwlb: 幂等落盘(按 url_hash 去重,重写 index 替代裸追加) - pipeline: xwlb 步骤传 --date {date} - run_extractor: xwlb publish_time 从 url 兜底解析真实播出日 - 新增 tests/test_run_xwlb.py 9 个测试 - 清理 60 天历史死数据(用户决策 A1) - docs: 文档重构收尾(删除 deployment.md,已在服务器直接修改) --- continuation.md | 35 +++++++++ scheduler/pipeline.py | 2 +- scripts/run_extractor.py | 21 +++++ scripts/run_xwlb.py | 103 +++++++++++++++++++----- tests/test_run_xwlb.py | 164 +++++++++++++++++++++++++++++++++++++++ 5 files changed, 306 insertions(+), 19 deletions(-) create mode 100644 tests/test_run_xwlb.py diff --git a/continuation.md b/continuation.md index 1e4d951..e6543c3 100644 --- a/continuation.md +++ b/continuation.md @@ -4,6 +4,41 @@ --- +## 本次完成 (2026-08-22) — xwlb 日期错位修复(方案 A) + +**背景问题**: +- P0-1: `run_xwlb` 硬编码抓前一天并落盘到「数据日」目录,而下游 extractor 按「处理日」扫目录 → 60 天联播数据从未进入知识库(实测: `data/processed/xwlb/` 仅 20260622 一天) +- P0-2: `run_xwlb` 用 `open("a")` 裸追加,同一天被多次调度触发导致 index.jsonl 重复行(实测 20260821: 70 行仅 24 唯一) + +**业务约束(保持不变)**: 当天日报需要前一天晚上(19:00 播出)的新闻联播内容。 + +**方案 A 实施**: + +1. `scripts/run_xwlb.py` 重构: + - 新增 `--date` 参数(处理日,默认今天),与管道其它步骤日期语义一致 + - 内部 `数据日 = 处理日 - 1`(业务约束:抓前一晚已播出的联播) + - **落盘目录改为处理日** `data/raw/xwlb/{处理日}/` → extractor/dedup/llm/embedding/qdrant 零改动即可处理 + - **幂等落盘**: 落盘前读已有 index.jsonl 构建 `url_hash` 集合,已存在条目跳过;index.jsonl 重写为去重完整集(原子写),不再裸追加 +2. `scheduler/pipeline.py`: xwlb 步骤命令加 `--date {date}` +3. `scripts/run_extractor.py`: xwlb 假 HTML 无时间节点 → 新增 `_fill_xwlb_publish_time`,从 url(`xwlb://YYYY-MM-DD/sid`)兜底解析真实播出日填充 `publish_time`(否则检索时间过滤/排序失效) +4. **历史死数据 A1 清理**: 删除 `data/raw/xwlb/` 62 个旧目录(20260622~20260822)+ `data/processed/xwlb/`(仅 20260622) + +**验证(生产端到端,20260822)**: +- 新增 `tests/test_run_xwlb.py` 9 个测试(处理日目录、幂等、空数据、坏日期),全部通过 +- 全量测试 **255 passed**(3 个 crawler 基线失败与本次无关) +- 实跑: 处理日=20260822 抓数据日 20260821 联播 23 条 → 二次运行新增 0/跳过 23(幂等)→ M2 提取 23/23 → M3 唯一 23 → M4 LLM 23/23 → M5 23/23 → M6 Qdrant 净增 23 +- 检索验证: `a-share search "国务院常务会议"` top2 命中 xwlb《李强主持召开国务院常务会议》,时间 2026-08-21(真实播出日)✓ +- 日报路径不受影响(仍走 `reporter._collect_xwlb` API 直读) + +**运维规则新增**: 每次代码升级完成后必须执行 `sudo systemctl restart a-share-research.service` + +**待办/遗留**: +- P0-3(crawler 无 --date,补跑历史日期静默空转)、P0-4(report 硬编码 date.today())、P1(dedup rc=1 静默/补跑含 cninfo/时区)未修,待用户决策 +- `run_extractor --force` 重提取会重复追加 `data/processed/{src}/{date}/index.jsonl`(仅影响该辅助索引,下游读 *.json 不受影响),暂不修 +- 历史 60 天联播未入库(按用户决策 A1 清理,不回补) + +--- + ## 本次完成 (2026-08-22) — 文档清理与重构 **目标:** 清理历史上多个 AI agent 文档残留,重构项目文档为 4 个核心文件。 diff --git a/scheduler/pipeline.py b/scheduler/pipeline.py index 993c54b..ea303b9 100644 --- a/scheduler/pipeline.py +++ b/scheduler/pipeline.py @@ -62,7 +62,7 @@ STEP_TIMEOUTS: dict[str, int] = { # crawler 不支持 --date,固定写当天目录; dedup 不加 --reset 以保持增量 STEP_COMMANDS: dict[str, list[str]] = { "crawler": ["uv", "run", "python", "-m", "scripts.run_crawler"], # 无 --date - "xwlb": ["uv", "run", "python", "-m", "scripts.run_xwlb"], + "xwlb": ["uv", "run", "python", "-m", "scripts.run_xwlb", "--date", "{date}"], # 处理日目录 "extractor": ["uv", "run", "python", "-m", "scripts.run_extractor", "--date", "{date}"], "dedup": ["uv", "run", "python", "-m", "scripts.run_dedup", "--date", "{date}"], "llm": ["uv", "run", "python", "-m", "scripts.run_event_extraction", "--date", "{date}"], diff --git a/scripts/run_extractor.py b/scripts/run_extractor.py index 135b21c..7d3aee9 100644 --- a/scripts/run_extractor.py +++ b/scripts/run_extractor.py @@ -115,9 +115,30 @@ def _process_one(rec: dict, raw_dir: Path, out_dir: Path, logger.warning("提取失败 {} {}: {}", rec["source_id"], rec["url"], e.reason) return None + # xwlb 的假 HTML 无时间节点,从 url(xwlb://YYYY-MM-DD/sid)兜底解析真实播出日 + if src_id == "xwlb" and article.publish_time is None: + article = _fill_xwlb_publish_time(article) + return _save_article(article, out_dir) +def _fill_xwlb_publish_time(article: Article) -> Article: + """从 xwlb url 解析播出日期填充 publish_time,失败时原样返回。""" + import re as _re + + m = _re.match(r"xwlb://(\d{4}-\d{2}-\d{2})", article.url) + if not m: + return article + try: + pt = datetime.strptime(m.group(1), "%Y-%m-%d") + except ValueError: + return article + return article.model_copy(update={ + "publish_time": pt, + "publish_time_raw": m.group(1), + }) + + def _process_cninfo(rec: dict, html_path: Path, out_dir: Path) -> Article | None: """处理 cninfo 公告记录:从 meta JSON 解析结构化数据。""" import re diff --git a/scripts/run_xwlb.py b/scripts/run_xwlb.py index 1a657ad..4f4b852 100644 --- a/scripts/run_xwlb.py +++ b/scripts/run_xwlb.py @@ -1,10 +1,24 @@ """M1 新闻联播 API 抓取脚本。 -从 doorcome API /api/xwlbFine/ 获取 AI 精编的新闻联播条目, -转换为与 Web 抓取兼容的格式(data/raw/xwlb/{date}/), -供 M2-M6 管道统一处理。 +《新闻联播》每天 19:00 播出:日报/知识库在早上生成时,当日报表只能引用 +前一天晚上已播出的联播,因此本脚本固定抓取「处理日 - 1 天」的节目, +但**落盘到「处理日」目录** data/raw/xwlb/{处理日}/。 -新闻联播晚间播出,始终抓取前一天数据,不依赖 --date 参数。 +日期语义说明: + - 处理日 = pipeline 本次批次的日期(--date,默认今天),与管道其它步骤一致; + - 数据日 = 处理日 - 1(昨晚 19:00 已播出的联播),业务约束来源: + 当天的日报需要前一天晚上的新闻联播内容; + - 目录名使用处理日,使 M2→M6(extractor/dedup/llm/embedding/qdrant) + 按同一日期扫目录即可处理 xwlb 数据,与其它新闻源完全一致。 + +幂等:同一天被定时任务多次触发(如 07:00/12:00/18:00/22:00)时, +数据日相同 → 抓取结果相同;按 url_hash 去重后只保留一份, +index.jsonl 重写为去重后的完整集(不再追加产生重复行)。 + +用法: + uv run python -m scripts.run_xwlb # 处理日 = 今天 + uv run python -m scripts.run_xwlb --date 20260823 # 指定处理日 + uv run python -m scripts.run_xwlb --output-root data/raw """ from __future__ import annotations @@ -33,19 +47,56 @@ def _url_hash(s: str) -> str: return hashlib.sha1(s.encode("utf-8")).hexdigest()[:16] +def _load_existing_index(index_path: Path) -> dict[str, str]: + """读取已有 index.jsonl,返回 {url_hash: 原始行},用于幂等去重。 + + 已存在同 url_hash 的条目视为已落盘,跳过写入;行内容原样保留。 + 文件不存在或损坏行忽略,不影响本次运行。 + """ + existing: dict[str, str] = {} + if not index_path.is_file(): + return existing + for line in index_path.read_text(encoding="utf-8").splitlines(): + line = line.strip() + if not line: + continue + try: + rec = json.loads(line) + except json.JSONDecodeError: + logger.warning("跳过 index.jsonl 非法行: {}", line[:80]) + continue + key = rec.get("url_hash") or rec.get("url") + if key: + existing[key] = line + return existing + + def main() -> int: parser = argparse.ArgumentParser(description="新闻联播 API 抓取 (M1-xwlb)") + parser.add_argument( + "--date", + default=date.today().strftime("%Y%m%d"), + help="处理日 YYYYMMDD(知识库批次日,默认今日);抓取该日前一天播出的联播", + ) parser.add_argument("--output-root", default="data/raw") parser.add_argument("--log-level", default="INFO") args = parser.parse_args() _setup_logger(args.log_level) - # 新闻联播晚间播出,始终抓取前一天 - day_str = (date.today() - timedelta(days=1)).strftime("%Y%m%d") - api_url = f"https://api.doorcome.cn/api/xwlbFine/?start_date={day_str}&end_date={day_str}" + # 处理日 → 数据日(前一晚已播出的联播) + try: + process_day = datetime.strptime(args.date, "%Y%m%d").date() + except ValueError: + logger.error("--date 格式错误: {!r},应为 YYYYMMDD", args.date) + return 2 + target_day = (process_day - timedelta(days=1)).strftime("%Y%m%d") + api_url = f"https://api.doorcome.cn/api/xwlbFine/?start_date={target_day}&end_date={target_day}" - logger.info("请求新闻联播 API: {}", api_url) + logger.info( + "请求新闻联播 API: {} (处理日={}, 数据日={})", + api_url, args.date, target_day, + ) try: req = urllib.request.Request(api_url) with urllib.request.urlopen(req, timeout=15) as resp: @@ -56,26 +107,34 @@ def main() -> int: raw_news = body.get("data", {}).get("news", []) if not raw_news: - logger.warning("{} 无新闻联播数据", day_str) + logger.warning("{} 无新闻联播数据", target_day) return 0 - # 准备输出目录 - out_dir = Path(args.output_root) / "xwlb" / day_str + # 输出目录: data/raw/xwlb/{处理日}/(目录 = pipeline 批次日) + out_dir = Path(args.output_root) / "xwlb" / args.date out_dir.mkdir(parents=True, exist_ok=True) - index_path = out_dir / "index.jsonl" + + # 幂等: 已有条目去重(同一天多次调度不会重复写入) + existing = _load_existing_index(index_path) saved = 0 + skipped = 0 + dup_rows = 0 for n in raw_news: sid = n.get("daily_sub_id", 0) title = n.get("news_title", "") content = n.get("news_improve", "") - news_day = n.get("news_days", day_str) + news_day = n.get("news_days", target_day) fake_url = f"xwlb://{news_day}/{sid}" h = _url_hash(fake_url) - html_file = f"{h}.html" + if h in existing: + skipped += 1 + continue # 幂等: 已落盘过,跳过写入 + + html_file = f"{h}.html" html_content = f""" {title}

{title}

{content}
@@ -95,13 +154,21 @@ def main() -> int: "url_hash": h, "html_file": html_file, } - with index_path.open("a", encoding="utf-8") as f: - f.write(json.dumps(meta, ensure_ascii=False) + "\n") + existing[h] = json.dumps(meta, ensure_ascii=False) saved += 1 - logger.info("新闻联播 {} 抓取完成: {} 条 -> {}", day_str, saved, out_dir) + # 重写 index.jsonl 为去重后的完整集(原子写,替代裸追加) + if saved or existing: + tmp = index_path.with_suffix(".jsonl.tmp") + tmp.write_text("\n".join(existing.values()) + "\n", encoding="utf-8") + tmp.replace(index_path) + + logger.info( + "新闻联播 数据日 {} 抓取完成: 新增 {} 条, 跳过已存在 {} 条 -> {}", + target_day, saved, skipped, out_dir, + ) return 0 if __name__ == "__main__": - raise SystemExit(main()) + raise SystemExit(main()) \ No newline at end of file diff --git a/tests/test_run_xwlb.py b/tests/test_run_xwlb.py new file mode 100644 index 0000000..a48ad70 --- /dev/null +++ b/tests/test_run_xwlb.py @@ -0,0 +1,164 @@ +"""run_xwlb 脚本测试(处理日目录语义 + 幂等去重)。 + +覆盖: + - 落盘目录 = 处理日(而非数据日): 处理日 20260823, 抓数据日 20260822 的联播, + 写入 data/raw/xwlb/20260823/ + - 处理日目录下 index.jsonl 的 source_id/url_hash 正确 + - 幂等: 同一处理日重复运行,index.jsonl 不产生重复行(url_hash 唯一) + - 日期格式错误返回 2,API 无数据返回 0 +""" + +from __future__ import annotations + +import json +from datetime import date, timedelta +from pathlib import Path +from unittest import mock + +from scripts.run_xwlb import _load_existing_index, _url_hash + + +# --------------------------------------------------------------------------- # +# 幂等辅助函数 +# --------------------------------------------------------------------------- # + +def test_url_hash_stable() -> None: + """url_hash 稳定且为 SHA1 前 16 位。""" + h1 = _url_hash("xwlb://2026-08-22/3") + h2 = _url_hash("xwlb://2026-08-22/3") + assert h1 == h2 + assert len(h1) == 16 + assert _url_hash("xwlb://2026-08-22/3") != _url_hash("xwlb://2026-08-22/4") + + +def test_load_existing_index_empty(tmp_path: Path) -> None: + """无 index.jsonl 时返回空 dict。""" + assert _load_existing_index(tmp_path / "index.jsonl") == {} + + +def test_load_existing_index_parses(tmp_path: Path) -> None: + """能读取已有行,key=url_hash;跳过损坏行。""" + p = tmp_path / "index.jsonl" + p.write_text('{"url_hash": "abc", "url": "x"}\nnot-json\n{"url_hash": "def", "url": "y"}\n', encoding="utf-8") + existing = _load_existing_index(p) + assert set(existing) == {"abc", "def"} + + +def test_load_existing_index_fallback_url(tmp_path: Path) -> None: + """无 url_hash 时用 url 兜底。""" + p = tmp_path / "index.jsonl" + p.write_text('{"url": "xwlb://2026-08-22/1"}\n', encoding="utf-8") + existing = _load_existing_index(p) + assert "xwlb://2026-08-22/1" in existing + + +# --------------------------------------------------------------------------- # +# 主流程(以 mock 网络方式调用 main) +# --------------------------------------------------------------------------- # + +def _make_api_body() -> dict: + """构造 doorcome xwlb API 响应体(3 条 + 1 条内容提要)。""" + news = [] + for sid in range(1, 5): + news.append({ + "daily_sub_id": sid, + "news_title": f"联播标题{sid}", + "news_improve": f"联播正文内容{sid}", + "news_days": "2026-08-22", + }) + return {"data": {"news": news}} + + +def _run_main(tmp_path: Path, process_day: str, body: dict | None = None) -> int: + """以 mock API 响应调用 scripts.run_xwlb.main。""" + import scripts.run_xwlb as mod + + body = body if body is not None else _make_api_body() + with mock.patch.object(mod.urllib.request, "urlopen") as m_urlopen: + resp = mock.MagicMock() + resp.read.return_value = json.dumps(body).encode("utf-8") + m_urlopen.return_value.__enter__.return_value = resp + # 直接调用 main 前的 argparse 不方便,这里调用内部逻辑等价于 main: + # 通过 subprocess 太重,改为构造 args 后调用主流程函数(main 内联,故复制流程)。 + # 简化:直接验证 _load_existing_index + 手写落盘逻辑的核心——用真实 main 需要 sys.argv。 + # 因此此处通过 mock sys.argv 调用 main。 + import sys + old_argv = sys.argv + sys.argv = [ + "run_xwlb", + "--date", process_day, + "--output-root", str(tmp_path), + "--log-level", "ERROR", + ] + try: + rc = mod.main() + finally: + sys.argv = old_argv + return rc + + +def test_main_writes_process_day_dir(tmp_path: Path) -> None: + """落盘目录 = 处理日(data/raw/xwlb/20260823),而非数据日 20260822。""" + rc = _run_main(tmp_path, "20260823") + assert rc == 0 + + out_dir = tmp_path / "xwlb" / "20260823" + assert out_dir.is_dir() + # 数据日目录不应存在 + assert not (tmp_path / "xwlb" / "20260822").is_dir() + + lines = (out_dir / "index.jsonl").read_text(encoding="utf-8").splitlines() + assert len(lines) == 4 # 4 条(含提要,但脚本不去重提要,只跳过 sid<=1 的逻辑在 reporter) + # 校验字段 + recs = [json.loads(l) for l in lines] + assert all(r["source_id"] == "xwlb" for r in recs) + assert all(r["stage"] == "article" for r in recs) + assert all(r["url_hash"] for r in recs) + # html 文件落盘 + html_files = list(out_dir.glob("*.html")) + assert len(html_files) == 4 + + +def test_main_idempotent_no_duplicate_rows(tmp_path: Path) -> None: + """同一处理日重复运行: index.jsonl 不产生重复行(url_hash 唯一)。""" + rc1 = _run_main(tmp_path, "20260823") + rc2 = _run_main(tmp_path, "20260823") + assert rc1 == 0 and rc2 == 0 + + lines = (tmp_path / "xwlb" / "20260823" / "index.jsonl").read_text(encoding="utf-8").splitlines() + hashes = [json.loads(l)["url_hash"] for l in lines] + assert len(hashes) == len(set(hashes)), "重复运行产生重复行" + assert len(lines) == 4, "第二次运行应全部跳过,行数不变" + + +def test_main_different_process_days_isolated(tmp_path: Path) -> None: + """不同处理日使用不同目录,互不污染。""" + _run_main(tmp_path, "20260823") + _run_main(tmp_path, "20260824") + assert (tmp_path / "xwlb" / "20260823").is_dir() + assert (tmp_path / "xwlb" / "20260824").is_dir() + # 20260824 的数据日是 20260823,若 mock 固定返回 2026-08-22 数据, + # 两个目录内容应相同字段结构,但互不影响 + n1 = len((tmp_path / "xwlb" / "20260823" / "index.jsonl").read_text(encoding="utf-8").splitlines()) + n2 = len((tmp_path / "xwlb" / "20260824" / "index.jsonl").read_text(encoding="utf-8").splitlines()) + assert n1 == 4 and n2 == 4 + + +def test_main_empty_news_returns_0(tmp_path: Path) -> None: + """API 无数据: 返回 0,不落盘。""" + rc = _run_main(tmp_path, "20260823", body={"data": {"news": []}}) + assert rc == 0 + assert not (tmp_path / "xwlb" / "20260823").is_dir() + + +def test_main_bad_date_returns_2(tmp_path: Path) -> None: + """--date 格式错误: 返回 2。""" + import sys + import scripts.run_xwlb as mod + old_argv = sys.argv + sys.argv = ["run_xwlb", "--date", "2026-13-99", "--output-root", str(tmp_path)] + try: + rc = mod.main() + finally: + sys.argv = old_argv + assert rc == 2 \ No newline at end of file