diff --git a/README.md b/README.md index ee3eee1..380d4d0 100644 --- a/README.md +++ b/README.md @@ -380,13 +380,38 @@ uv run python -m scripts.run_scheduler --once --date 20260616 # 只执行部分步骤(逗号分隔) uv run python -m scripts.run_scheduler --once --steps crawler,extractor,llm +# 断点续跑:跳过连续成功步骤,从上次失败/未执行步骤继续 +uv run python -m scripts.run_scheduler --once --resume + # 启动定时守护进程(按 .env 中 SCHEDULE_TIMES 自动触发) uv run python -m scripts.run_scheduler ``` 定时时间由 `.env` 中 `SCHEDULE_TIMES` 控制(默认 `07:00,12:00,18:00,22:00`)。 -Pipeline 总耗时约 4-5 分钟(100 篇文章),其中 M1 抓取(含浏览器渲染)最耗时(~3 分钟)。 +### 增量处理与断点续跑 + +各步骤默认「产物存在即跳过」,中断/重跑不会重复处理已完成部分(API 密集的 +M4/M5 不会重复计费);需要全量重跑时加 `--force`。 + +| 步骤 | 增量机制 | 全量重跑 | +| --- | --- | --- | +| M1 crawler | `seen_urls.txt` 记录历史 URL,列表页链接按 hash 过滤 | 清 `data/raw/*/seen_urls.txt` | +| M2 extractor | 输出目录已有 `{url_hash}.json` 即跳过提取 | `--force` | +| M3 dedup | 指纹库增量判重(重复不入库) | `--reset` | +| M4 llm | `data/events/{day}/{url_hash}.json` 已存在即跳过(不重复调 API) | `--force` | +| M5 embedding | `data/embeddings/{day}/{url_hash}.json` 已存在即跳过 | `--force` | +| M6 qdrant | upsert 幂等(url_hash 为 point ID) | `--recreate` | + +**断点续跑** `pipeline --once --resume`(仅全链路,不能与 `--steps` 同用): + +- 每次运行把各步骤结果(ok/failed + 退出码 + 耗时)写入 `data/pipeline/state.json`(按日期隔离); +- `--resume` 读取该日期的状态,跳过连续成功的步骤,从第一个失败/未执行步骤继续执行到结尾; +- 中断(人为 Ctrl+C / 报错 / 超时)后重新执行同一条命令即可自动从断点继续; +- 定时任务(守护模式)始终全量执行,不受 resume 影响。 + +**实测效果**(20260616 数据,100 篇全链路):增量重跑时 M2 从全量重新提取降至约 2.5s(126 篇跳过), +M4 减少 100 次 LLM API 调用、M5 减少 100 次 embedding API 调用,均为秒级完成。 日志:`logs/scheduler.log` diff --git a/continuation.md b/continuation.md index 8f5dd5c..bf1baed 100644 --- a/continuation.md +++ b/continuation.md @@ -1,6 +1,6 @@ # continuation.md -> `checkpoint` @ 2026-08-11 17:00 +> `checkpoint` @ 2026-08-12 10:30 --- @@ -15,44 +15,40 @@ | 调度器 | APScheduler,systemd `a-share-research.service`(pi5);每天 07:00 首次任务生成日报(12/18/22 点不生成) | | LLM | 场景化配置 `configs/llm_models.yaml`(4 场景: event_extraction/daily_report/stock_report/embedding);YAML 优先、`.env` 兜底;模型必须显式配置,无内置兜底 | | 去重 | 多源记录:指纹库 `source_ids` 列 + uniques JSON `sources` 字段 + `data/deduped/{day}/sources.json` | +| 增量/断点 | M2/M4/M5 产物存在即跳过(`--force` 全量);`pipeline --once --resume` 断点续跑(状态 `data/pipeline/state.json`) | | 服务器 | `pi@192.168.1.160`(生产)/ `pi@192.168.1.10`(DB 隧道宿主) | | 抓取方式 | js_render=false → httpx 直连;js_render=true → Playwright | --- -## 本次完成 (2026-08-11) — 大模型场景化配置 + 去重多源记录 +## 本次完成 (2026-08-12) — 增量处理与 pipeline 断点续跑 -**目标:** ① 梳理全部 AI 大模型使用点,新增 `configs/llm_models.yaml` 按场景独立配置 provider/model;② 去重时记录一条唯一新闻的全部来源。 +**目标:** ① 全链路中断后可从断点恢复;② 各子任务排除已处理文件,避免全量重跑与重复 API 计费。 -**1. 大模型使用点梳理(共 4 个场景,详见 configs/llm_models.yaml 内注释):** -- `event_extraction`(M4 投资事件抽取,JSON mode,llm/extractor.py) -- `daily_report`(日报 AI 摘要,scheduler/reporter.py) -- `stock_report`(个股 AI 要点分析,scheduler/stock_reporter.py) -- `embedding`(向量化,dashscope 远程 / local-bge 本地,embedding/remote.py + local.py) -- 非使用点确认:crawler 纯抓取、MCP 仅复用 embedding、run_xwlb 抓外部「AI 精编」数据源 +**1. 各步骤增量处理(产物存在即跳过,`--force` 全量):** +- M2 `run_extractor.py`:输出目录已有 `{url_hash}.json` 即跳过提取,仅回补 index 行;`--force` 重建;成功率统计含跳过项(修复全跳过时误报 rc=1) +- M4 `run_event_extraction.py`:`data/events/{day}/{url_hash}.json` 已存在即跳过(**不重复调用 LLM API**);`--force` 全量;failed.jsonl 只保留本次失败、index 累积追加 +- M5 `run_embedding.py`:`data/embeddings/{day}/{url_hash}.json` 已存在即跳过(**不重复调用 embed API**);`--force` 全量 +- M1(seen_urls 增量)/ M3(指纹库判重)/ M6(upsert 幂等)为既有能力,README 汇总成表 -**2. 场景配置实现(优先级: CLI 显式参数 > YAML > .env > 内置默认):** -- 新增 `configs/loader.py`(lru_cache 读 llm_models.yaml)+ `configs/__init__.py` -- `llm/client.py`:`load_llm_config(scene=...)` 支持场景;`LLMConfig` 增加 `max_attempts`;模型缺失仍报错(无内置兜底) -- `llm/extractor.py`:`extract_event(_async)` 的 max_attempts 默认取 `config.max_attempts` -- `embedding/factory.py` + `remote.py` + `local.py`:provider/model/api_key_env/base_url_env/batch_limit 支持场景覆盖 -- `scheduler/reporter.py`(daily_report)+ `stock_reporter.py`(stock_report):接入场景,temperature 取配置 -- YAML 中 provider/model 默认留空 → 回退 .env,**现有部署零改动兼容** +**2. pipeline 断点续跑(scheduler/pipeline.py + run_scheduler.py):** +- 新增 `data/pipeline/state.json`(按日期隔离,记录每步骤 ok/failed + 退出码 + 耗时),原子写 +- `run_pipeline(resume=True)` 跳过连续成功前缀,从首个失败/未执行步骤继续执行到结尾 +- `pipeline --once --resume`(默认全量不变;`--resume` 与 `--steps` 互斥报错);定时守护模式不受影响 -**3. 去重多源记录:** -- `dedup/models.py`:`Fingerprint.source_ids`(validator 保主源居首+去重);`DedupResult` 增 `matched_source_id`/`all_source_ids` -- `dedup/store.py`:指纹库加 `source_ids` 列,旧库自动 ALTER 迁移,旧数据回退 `[source_id]` -- `dedup/deduper.py`:`ingest` 命中重复时把新源合并进匹配指纹 -- `scripts/run_dedup.py`:uniques JSON 附加 `sources` 字段;重复命中时仅更新 sources 不覆盖原文;输出 `data/deduped/{day}/sources.json` 汇总 +**验证(Mac 本地,20260616 数据 100 篇):** +- pytest **225 passed**(新增 tests/test_incremental.py 10 个:M2/M4/M5 跳过、状态记录、resume 续跑、resume 全完成 noop、--resume+--steps 互斥);crawler 3 个基线失败仍与本次无关 +- ruff 零新增(9 个基线错误不变) +- 端到端:M2 增量重跑 140 条跳过 126 条,2.5s 完成、rc=0(修复前误报失败);state.json 正确记录 extractor ok -**验证(Mac 本地):** -- 全量 pytest:**215 passed**(仅 crawler 3 个 retry mock 失败为基线预存在问题,与本改动无关) -- 端到端人工构造 3 源同文:1 条唯一 + sources.json `["cls","eastmoney","sina"]` + 指纹库 source_ids 列正确 -- ruff:9 个错误均为基线既有(crawler/cninfo.py 未用 import、reporter.py L5/L4 命名),本次零新增 +**效果评估(100 篇规模中断重跑场景):** +- M4 减少约 100 次 LLM API 调用、M5 减少约 100 次 embedding API 调用 → 中断恢复不再重复计费,耗时从分钟级降至秒级 +- M2 重跑从全量 GNE 提取(分钟级)降至约 2.5s +- 断点恢复操作:中断后直接重跑同一条 `pipeline --once --resume` 命令即可 **待办:** -- 生产同步:代码 + `configs/llm_models.yaml` scp 到 pi5(注意 rsync 排除规则含 configs/*.yaml,需显式同步),重启 `a-share-research` 生效 -- 首次同步前 pi5 无 YAML → 全部回退 .env,行为不变,可平滑切换 +- 同步 pi5(代码 + 文档),重启 `a-share-research`;首次同步后 pi5 的 `data/pipeline/state.json` 不存在 → resume 按全量处理,行为安全 +- git 提交(本次改动尚未提交) ---