From 7b33182f6796024e9eb2bb2e66cf5c3660f35d86 Mon Sep 17 00:00:00 2001 From: simon Date: Sun, 23 Aug 2026 05:40:52 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E8=A1=A5=E8=B7=91?= =?UTF-8?q?=E8=AF=AF=E5=90=AB=20cninfo=20=E4=B8=89=E6=AD=A5=E4=B8=8E?= =?UTF-8?q?=E6=97=B6=E5=8C=BA=E4=B8=8D=E4=B8=80=E8=87=B4=20(P1-2/P1-3)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - P1-2: 提取 DEFAULT_NEWS_STEPS(全链路去 report+cninfo),定时/补跑/默认三处统一; 启动补跑不再误执行 cninfo 公告管道 - P1-3: 新增 scheduler/timeutil.py(调度时区统一入口,SCHEDULE_TZ 可覆盖, 默认 Asia/Shanghai);run_scheduler 定时/补跑/--once、crawler 补跑保护、 reporter 兜底日期全部改用调度时区,避免系统时区非上海时日期错位 - 新增 6 个测试;全量 282 passed --- continuation.md | 18 +++++++ scheduler/__init__.py | 2 + scheduler/pipeline.py | 17 ++++-- scheduler/reporter.py | 5 +- scheduler/timeutil.py | 36 +++++++++++++ scripts/run_scheduler.py | 32 +++++++----- tests/test_incremental.py | 107 +++++++++++++++++++++++++++++++++++++- 7 files changed, 197 insertions(+), 20 deletions(-) create mode 100644 scheduler/timeutil.py diff --git a/continuation.md b/continuation.md index a0b4ecf..b98dd1b 100644 --- a/continuation.md +++ b/continuation.md @@ -4,6 +4,24 @@ --- +## 本次完成 (2026-08-23) — P1-2/P1-3 补跑步骤与时区一致性修复 + +**P1-2(守护进程补跑误含 cninfo 三步)**: +- 原因:`scripts/run_scheduler.py` 启动补跑逻辑 `steps = [k for k in STEP_COMMANDS if k != "report"]` 只排除 report,**漏掉 cninfo_crawl/cninfo_extract/cninfo_pdf**;错过定时任务重启补跑时会额外执行整套公告管道(且 cninfo_crawl 无 --date,补跑历史日期静默空转) +- 修复:提取共享常量 `scheduler.pipeline.DEFAULT_NEWS_STEPS`(= 全链路去除 report 与 cninfo 三步),定时任务、补跑、run_pipeline 默认三处统一引用 + +**P1-3(调度时区与 date.today() 不一致)**: +- 原因:cron 触发器显式用 `Asia/Shanghai`,但 `date.today()`/`datetime.now()` 取**系统时区**;若系统时区非上海(如容器 UTC),07:00 上海(=前一日 23:00 UTC)触发时日期会错一天,整条 pipeline 落错日目录 +- 修复:新增 `scheduler/timeutil.py`(schedule_tz/today_str/now,env `SCHEDULE_TZ` 可覆盖,默认 Asia/Shanghai);`run_scheduler` 定时/补跑/`--once`、`pipeline.run_step` crawler 补跑保护、`reporter.generate_report` 兜底日期全部改用它 + +**测试**: +- 新增 6 个(DEFAULT_NEWS_STEPS 不含 cninfo/run_pipeline 默认步骤行为/补跑源码引用/时区两例/--once 缺省日期),并入 `tests/test_incremental.py` +- 全量 **282 passed**,3 个 crawler 基线失败与本次无关;ruff 干净(5 个 reporter 既有问题未动) + +**验证**:DEFAULT_NEWS_STEPS = [crawler, xwlb, extractor, dedup, llm, embedding, qdrant] ✓;today_str = 20260823(Asia/Shanghai)✓ + +--- + ## 本次完成 (2026-08-22) — 多源新闻记录链路修复(方案 A + B) **用户需求**:一条新闻有多个来源时,全部来源都要记录并可见;此前"找不到多源"。 diff --git a/scheduler/__init__.py b/scheduler/__init__.py index 7fd674b..430db13 100644 --- a/scheduler/__init__.py +++ b/scheduler/__init__.py @@ -7,6 +7,7 @@ """ from .pipeline import ( + DEFAULT_NEWS_STEPS, STEP_COMMANDS, STEP_TIMEOUTS, PipelineResult, @@ -16,6 +17,7 @@ from .pipeline import ( ) __all__ = [ + "DEFAULT_NEWS_STEPS", "STEP_COMMANDS", "STEP_TIMEOUTS", "PipelineResult", diff --git a/scheduler/pipeline.py b/scheduler/pipeline.py index a025162..083b6c2 100644 --- a/scheduler/pipeline.py +++ b/scheduler/pipeline.py @@ -16,11 +16,13 @@ import os import subprocess import time from dataclasses import dataclass, field -from datetime import date, datetime +from datetime import datetime from pathlib import Path from loguru import logger +from .timeutil import today_str + # 断点状态文件(按日期隔离,记录每步骤结果) DEFAULT_STATE_PATH = Path("data/pipeline/state.json") @@ -74,6 +76,14 @@ STEP_COMMANDS: dict[str, list[str]] = { "cninfo_pdf": ["uv", "run", "a-share", "cninfo", "--enrich-pdf", "--pdf-limit", "100"], } +# 新闻链路默认步骤:全链路除去 report(单独追加)与 cninfo 独立管道(P1-2: +# 补跑/定时任务若误含 cninfo 三步,会重复执行整套公告管道,且 cninfo_crawl +# 无 --date 参数,补跑历史日期时会静默空转)。 +DEFAULT_NEWS_STEPS: list[str] = [ + k for k in STEP_COMMANDS + if k not in ("report", "cninfo_crawl", "cninfo_extract", "cninfo_pdf") +] + @dataclass class StepResult: @@ -223,7 +233,8 @@ def run_step(name: str, date_str: str) -> StepResult: # 补跑保护:抓取类步骤只能产生当天数据(网站首页只含当前内容,历史文章已滚走), # 补跑历史日期时跳过抓取并告警,复用已有 data/raw/*/{date_str} 数据。 # 注意:仅保护 crawler; xwlb 走 API 支持任意历史日期,无需跳过。 - if name == "crawler" and date_str != date.today().strftime("%Y%m%d"): + # "今天"按调度时区判定(P1-3),避免系统时区与 cron 时区不一致时错判。 + if name == "crawler" and date_str != today_str(): started = datetime.now() logger.warning( "补跑模式:首页只含当天内容,无法补抓 {};跳过抓取,复用已有 data/raw/*/{}", @@ -330,7 +341,7 @@ def run_pipeline( 记录,跳过连续成功的步骤,从第一个失败/未执行步骤继续。 state_path: 断点状态文件路径(测试可注入)。 """ - names = steps or [k for k in STEP_COMMANDS if k not in ("report", "cninfo_crawl", "cninfo_extract", "cninfo_pdf")] + names = steps or list(DEFAULT_NEWS_STEPS) result = PipelineResult(started_at=datetime.now()) state = _load_pipeline_state(state_path) diff --git a/scheduler/reporter.py b/scheduler/reporter.py index abe9cd7..70d5bb6 100644 --- a/scheduler/reporter.py +++ b/scheduler/reporter.py @@ -1056,7 +1056,10 @@ def generate_report(day_str: str | None = None, *, upload: bool = True) -> int | `upload` 参数保留以兼容 scheduler/pipeline.py 调用,已无实际作用。 返回 report_id(成功)或 None(无数据/失败)。 """ - day_str = day_str or date.today().strftime("%Y%m%d") + # 兜底日期按调度时区取(P1-3),与 pipeline 传入的 date_str 语义一致 + if not day_str: + from .timeutil import today_str + day_str = today_str() logger.info("生成日报: {}", day_str) # 收集数据 diff --git a/scheduler/timeutil.py b/scheduler/timeutil.py new file mode 100644 index 0000000..0803f23 --- /dev/null +++ b/scheduler/timeutil.py @@ -0,0 +1,36 @@ +"""调度时区工具。 + +背景(P1-3): + 定时任务 cron 触发时区与"今天"的判定必须一致。此前 cron 用 + Asia/Shanghai,而 `date.today()` / `datetime.now()` 取系统时区—— + 若系统时区不是 Asia/Shanghai(如 UTC),07:00 上海(=前一日 23:00 UTC) + 触发时 `date.today()` 会返回错误日期,整条 pipeline 落在错日目录。 + +统一入口: + schedule_tz(): 调度时区(env SCHEDULE_TZ 可覆盖,默认 Asia/Shanghai) + today_str(): 按调度时区返回 YYYYMMDD +""" + +from __future__ import annotations + +import os +from datetime import datetime +from zoneinfo import ZoneInfo + +# 默认调度时区(与 .env 的 SCHEDULE_TIMES 语义一致) +DEFAULT_SCHEDULE_TZ = "Asia/Shanghai" + + +def schedule_tz() -> ZoneInfo: + """调度时区(env SCHEDULE_TZ 可覆盖,默认 Asia/Shanghai)。""" + return ZoneInfo(os.environ.get("SCHEDULE_TZ") or DEFAULT_SCHEDULE_TZ) + + +def now() -> datetime: + """当前时刻(调度时区,aware)。""" + return datetime.now(schedule_tz()) + + +def today_str() -> str: + """按调度时区返回今天的 YYYYMMDD。""" + return now().strftime("%Y%m%d") diff --git a/scripts/run_scheduler.py b/scripts/run_scheduler.py index 5d11f4f..29b99ae 100644 --- a/scripts/run_scheduler.py +++ b/scripts/run_scheduler.py @@ -16,15 +16,16 @@ from __future__ import annotations import argparse import signal import sys -from datetime import date, datetime from pathlib import Path from typing import Any from dotenv import load_dotenv from loguru import logger -from scheduler import STEP_COMMANDS, run_pipeline +from scheduler import DEFAULT_NEWS_STEPS, run_pipeline from scheduler.stock_reporter import generate_all_stock_reports +from scheduler.timeutil import now as tz_now +from scheduler.timeutil import schedule_tz, today_str def _setup_logger(level: str) -> None: @@ -69,7 +70,7 @@ def _once(args: argparse.Namespace) -> int: if args.resume and args.steps: logger.error("--resume 与 --steps 不能同时使用(断点续跑针对全链路)") return 2 - run_pipeline(args.date, steps=steps, resume=args.resume) + run_pipeline(args.date or today_str(), steps=steps, resume=args.resume) return 0 @@ -92,17 +93,18 @@ def _daemon(args: argparse.Namespace) -> int: scheduler = BackgroundScheduler() - # 包装函数:每次触发时重新计算日期,避免 date.today() 在注册时冻结。 + # 包装函数:每次触发时按调度时区重新计算日期(P1-3), + # 避免 date.today() 在注册时冻结或与 cron 时区不一致。 def _scheduled_pipeline(steps: list[str] | None = None) -> None: - run_pipeline(date.today().strftime("%Y%m%d"), steps=steps) + run_pipeline(today_str(), steps=steps) for hour, minute in times: - trigger = CronTrigger(hour=hour, minute=minute, timezone="Asia/Shanghai") + trigger = CronTrigger(hour=hour, minute=minute, timezone=str(schedule_tz())) is_first = (hour == first_hour and minute == first_minute) job_kwargs: dict | None = None if is_first: job_kwargs = { - "steps": [k for k in STEP_COMMANDS if k not in ("report", "cninfo_crawl", "cninfo_extract", "cninfo_pdf")] + ["report"] + "steps": list(DEFAULT_NEWS_STEPS) + ["report"] } scheduler.add_job( _scheduled_pipeline, @@ -118,7 +120,7 @@ def _daemon(args: argparse.Namespace) -> int: cninfo_raw = os.environ.get("CNINFO_SCHEDULE_TIME", "06:30") cninfo_parts = cninfo_raw.split(":") cninfo_h, cninfo_m = int(cninfo_parts[0]), int(cninfo_parts[1]) if len(cninfo_parts) > 1 else 0 - cninfo_trigger = CronTrigger(hour=cninfo_h, minute=cninfo_m, timezone="Asia/Shanghai") + cninfo_trigger = CronTrigger(hour=cninfo_h, minute=cninfo_m, timezone=str(schedule_tz())) cninfo_steps = ["cninfo_crawl", "cninfo_extract", "cninfo_pdf", "dedup", "llm", "embedding", "qdrant"] scheduler.add_job( @@ -135,7 +137,7 @@ def _daemon(args: argparse.Namespace) -> int: if stock_raw: stock_parts = stock_raw.split(":") stock_h, stock_m = int(stock_parts[0]), int(stock_parts[1]) if len(stock_parts) > 1 else 0 - stock_trigger = CronTrigger(hour=stock_h, minute=stock_m, timezone="Asia/Shanghai") + stock_trigger = CronTrigger(hour=stock_h, minute=stock_m, timezone=str(schedule_tz())) scheduler.add_job( generate_all_stock_reports, trigger=stock_trigger, @@ -159,7 +161,7 @@ def _daemon(args: argparse.Namespace) -> int: logger.info("调度器已启动,等待触发... (按 Ctrl+C 退出)") # 启动时检查是否有因重启/宕机错过的定时任务,30 分钟内补跑 - now = datetime.now() + now = tz_now() for hour, minute in times: scheduled = now.replace(hour=hour, minute=minute, second=0, microsecond=0) missed_minutes = (now - scheduled).total_seconds() / 60 @@ -168,10 +170,12 @@ def _daemon(args: argparse.Namespace) -> int: "检测到错过的定时任务 {:02d}:{:02d} ({} 分钟前),立即补跑一次", hour, minute, int(missed_minutes), ) - steps = [k for k in STEP_COMMANDS if k != "report"] + # P1-2:补跑只跑新闻链路(DEFAULT_NEWS_STEPS 不含 report 与 cninfo 三步), + # cninfo 公告管道由其自身定时任务负责,避免重复执行整套公告管道。 + steps = list(DEFAULT_NEWS_STEPS) if (hour, minute) == sorted_times[0]: steps.append("report") - run_pipeline(date.today().strftime("%Y%m%d"), steps=steps) + run_pipeline(today_str(), steps=steps) import contextlib @@ -186,8 +190,8 @@ def main() -> int: parser = argparse.ArgumentParser(description="A 股新闻定时任务 (M7)") parser.add_argument("--once", action="store_true", help="立即执行一次全链路") parser.add_argument( - "--date", default=date.today().strftime("%Y%m%d"), - help="日期 YYYYMMDD (仅 --once 模式)", + "--date", default=None, + help="日期 YYYYMMDD (仅 --once 模式,默认按调度时区取今天)", ) parser.add_argument("--steps", default=None, help="仅执行指定步骤,逗号分隔 (如 crawler,extractor)") diff --git a/tests/test_incremental.py b/tests/test_incremental.py index 6b0eda3..e0e0ad2 100644 --- a/tests/test_incremental.py +++ b/tests/test_incremental.py @@ -8,6 +8,7 @@ from __future__ import annotations import json +from datetime import UTC from pathlib import Path from types import SimpleNamespace @@ -306,8 +307,7 @@ def test_crawler_today_runs_normally(monkeypatch: pytest.MonkeyPatch) -> None: def test_report_step_uses_date_str(monkeypatch: pytest.MonkeyPatch) -> None: """report 步骤必须使用传入的 date_str,而非 date.today()(P0-4)。""" - from scheduler import pipeline - from scheduler import reporter + from scheduler import pipeline, reporter received: list[str] = [] @@ -346,3 +346,106 @@ def test_pipeline_backfill_skips_crawler_keeps_rest( assert all(s.success for s in result.steps) state = pipeline._load_pipeline_state(state_path) assert state["20260101"]["crawler"]["status"] == "ok" + + +# --------------------------------------------------------------------------- # +# P1-2: 默认步骤不含 cninfo 独立管道 / P1-3: 调度时区统一 +# --------------------------------------------------------------------------- # + +def test_default_news_steps_excludes_cninfo_and_report() -> None: + """新闻链路默认步骤不含 report 与 cninfo 三步(P1-2)。""" + from scheduler import DEFAULT_NEWS_STEPS + + assert "report" not in DEFAULT_NEWS_STEPS + assert "cninfo_crawl" not in DEFAULT_NEWS_STEPS + assert "cninfo_extract" not in DEFAULT_NEWS_STEPS + assert "cninfo_pdf" not in DEFAULT_NEWS_STEPS + # 新闻链路核心步骤齐全 + assert {"crawler", "xwlb", "extractor", "dedup", "llm", + "embedding", "qdrant"} <= set(DEFAULT_NEWS_STEPS) + + +def test_run_pipeline_default_steps_excludes_cninfo( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + """run_pipeline 默认 steps 不含 cninfo(只跑新闻链路)。""" + from scheduler import pipeline + from scheduler.timeutil import today_str + + calls: list[str] = [] + + def _fake_run(cmd, timeout=None): # noqa: ARG001 + name = next(c.split(".")[-1] for c in cmd if "scripts.run_" in c) + calls.append(name) + return SimpleNamespace(returncode=0) + + monkeypatch.setattr(pipeline.subprocess, "run", _fake_run) + state_path = tmp_path / "state.json" + today = today_str() # 今天:避免 crawler 补跑保护跳过 + result = pipeline.run_pipeline( + today, state_path=state_path, + ) + assert result.all_success is True + assert calls == ["run_crawler", "run_xwlb", "run_extractor", "run_dedup", + "run_event_extraction", "run_embedding", "run_qdrant_ingest"] + assert "cninfo" not in " ".join(calls) + + +def test_scheduler_daemon_backfill_uses_default_news_steps( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """守护进程补跑 steps 使用 DEFAULT_NEWS_STEPS(不含 cninfo 三步)。""" + import inspect + + import scripts.run_scheduler as rs + from scheduler import DEFAULT_NEWS_STEPS + + # 验证补跑逻辑引用的常量:通过源码断言 + 常量内容双重保证 + src = inspect.getsource(rs._daemon) + assert "list(DEFAULT_NEWS_STEPS)" in src + assert "cninfo_crawl" not in DEFAULT_NEWS_STEPS + + +def test_today_str_matches_schedule_timezone(monkeypatch: pytest.MonkeyPatch) -> None: + """today_str 按调度时区返回 YYYYMMDD(P1-3)。""" + from datetime import datetime + from zoneinfo import ZoneInfo + + from scheduler.timeutil import today_str + + s = today_str() + assert len(s) == 8 and s.isdigit() + # 与调度时区(默认 Asia/Shanghai)当前日期一致 + expect = datetime.now(ZoneInfo("Asia/Shanghai")).strftime("%Y%m%d") + assert s == expect + + +def test_today_str_respects_env_override(monkeypatch: pytest.MonkeyPatch) -> None: + """SCHEDULE_TZ 环境变量可覆盖调度时区(P1-3)。""" + from datetime import datetime + + from scheduler.timeutil import today_str + + # 覆盖为 UTC 后,today_str 应返回 UTC 日期(而非默认 Asia/Shanghai) + monkeypatch.setenv("SCHEDULE_TZ", "UTC") + s = today_str() + expect = datetime.now(UTC).strftime("%Y%m%d") + assert s == expect + + +def test_once_uses_today_str_when_date_missing(monkeypatch: pytest.MonkeyPatch) -> None: + """--once 不带 --date 时按调度时区取今天(P1-3)。""" + import scripts.run_scheduler as rs + + captured: dict[str, str] = {} + + def _fake_pipeline(date_str, steps=None, resume=False): # noqa: ARG001 + captured["date"] = date_str + + monkeypatch.setattr(rs, "run_pipeline", _fake_pipeline) + from types import SimpleNamespace + + args = SimpleNamespace(steps=None, resume=False, date=None) + rc = rs._once(args) + assert rc == 0 + assert captured["date"] == rs.today_str()