Author SHA1 Message Date
simon 951276a313 feat: 日报多来源入库(news_event.sources 列,JSON 数组,主源居首)
- EventRow 新增 sources 字段;init_schema 幂等 ALTER 补列(MariaDB)
- save_report 写入 sources JSON;_article_sources_full 完整来源列表
- source 字段保持 ≤3 拼接兼容;新增多源完整列表测试
- 已部署 pi5:sources 列写入验证通过(单源数组);完整多源待新数据
2026-08-12 13:06:01 +08:00
simon 32e99aeaf6 chore: .gitignore 忽略 data/dedup/ 指纹库(运行产物) 2026-08-12 11:23:35 +08:00
simon 012cb615d3 feat: 脚本执行显性输出当前阶段与 AI 模型供应商/名称
- pipeline.sh 每阶段 banner(━━━ M4 翻译+事件抽取(AI 大模型: deepseek / deepseek-v4-flash)━━━)
- ai_model_info() 从 system.yaml 读取场景模型(translation/daily_report/embedding)
- Python 层:M4/M5/日报 显性打印 provider/model(场景标注)
- 已同步 pi5 实测:客户端初始化日志含供应商/模型
2026-08-12 11:20:51 +08:00
simon bdd936a3e4 feat: 全流程终端实时进度输出 + 步骤日志
- step_run 增加耗时统计(▶ 开始 / ✔ 完成 Ns / ✗ 失败 exit=N)与 STEP_LOG 支持
- domestic_full.sh / pipeline.sh 去掉 tail 截断,输出实时透传终端并 tee 到 logs/*.log
- 退出码用 PIPESTATUS[0] 捕获(规避 pipefail/tee 吞码);全角标点后变量用花括号包裹
- 已同步 pi5 实测:M3 31s / M6 20s 耗时显示正常,日志文件生成
2026-08-12 11:18:08 +08:00
simon 06b00b7f49 feat: 全流程脚本 --resume 断点续跑 + 各步骤文件级增量说明
- 新增 scripts/_step_state.sh:步骤状态库(data/run_state/{date}.state,按天换新)
- pipeline.sh / domestic_full.sh 支持 --resume,按步骤跳过已完成项
- step_run 显式检查退出码(规避 bash set -e 条件上下文陷阱),失败步骤不标记
- M1 保持部分源失败容忍(step_should_run/step_mark 手动控制)
- 确认并文档化各步骤文件级增量:M2/M4/M5 按文件存在跳过、M3 指纹库、M6 upsert 幂等
- 已同步 pi5 并验证(--badarg exit=1、step_run 机制正常)
2026-08-12 11:06:29 +08:00
simon c72a5ed13a feat: AI 模型按场景独立配置(llm_scenes)+ 去重多来源合并
- system.yaml 新增 llm_scenes(translation / daily_report,含用途/方法/模型要求说明)
- load_llm_config(scene=...) 场景覆盖;日报摘要 temperature 0.3 硬编码 → 配置
- M3 去重:唯一篇记录 source_ids(跨源重复合并,首个来源为 source_id)
- source_ids 经翻译透传至 events,日报事件 source 多来源拼接展示(≤3 个)
- 已部署 pi5:merge 实证 10 源合并;日报 report_id=224 正常入库
2026-08-12 09:56:55 +08:00
simon 9aa44e610d feat: LLM/Embedding 重试次数纳入 system.yaml 统一配置
- system.yaml: llm.max_retries(死配置) 改为 llm.max_attempts;embedding 段新增 max_attempts
- llm/client.py: LLMConfig 新增 max_attempts(默认 3 兜底),load_llm_config 读取配置
- llm/extractor.py: translate_and_extract(_async) max_attempts=None 时取 config.max_attempts
- embedding/client.py: load_embedding_config 读取 embedding.max_attempts
- scheduler/reporter.py: _call_llm_simple max_retries=None 时取 max_attempts-1
- 新增 2 个配置读取测试;已同步 pi5 验证
2026-08-05 08:48:43 +08:00
23 changed files with 711 additions and 55 deletions
+2
View File
@@ -27,10 +27,12 @@ Thumbs.db
data/raw/ data/raw/
data/processed/ data/processed/
data/deduped/ data/deduped/
data/dedup/
data/events/ data/events/
data/embeddings/ data/embeddings/
data/qdrant_storage/ data/qdrant_storage/
data/reports/ data/reports/
data/run_state/
# Logs # Logs
logs/*.log logs/*.log
+45 -1
View File
@@ -104,6 +104,12 @@ sudo systemctl enable --now privoxy
# 全流程:M1 抓取 → M2→M6 管道 → 日报 # 全流程:M1 抓取 → M2→M6 管道 → 日报
# 每天 06:00 / 12:00 / 18:00 / 22:00 自动执行 # 每天 06:00 / 12:00 / 18:00 / 22:00 自动执行
bash scripts/domestic_full.sh bash scripts/domestic_full.sh
# 中断恢复:跳过当天已完成步骤(状态存 data/run_state/{date}.state)
# 支持步骤级断点:M1_crawl / M2_extract / M3_dedup / M4_translate / M5_embed / M6_index / report
bash scripts/domestic_full.sh --resume
# 仅 M2→M6 管道(同样支持 --resume)
bash scripts/pipeline.sh --resume
``` ```
### 单步执行 ### 单步执行
@@ -140,11 +146,49 @@ ls -lt data/reports/
| 文件 | 用途 | | 文件 | 用途 |
|------|------| |------|------|
| `configs/sources.yaml` | 英文财经新闻源定义(13 个源) | | `configs/sources.yaml` | 英文财经新闻源定义(13 个源) |
| `configs/system.yaml` | 系统级业务参数(超时/并发/LLM 模型等) | | `configs/system.yaml` | 模型按场景(`llm_scenes`)/重试/阈值等系统配置 |
| `configs/profiles/8g_headful.yaml` | Pi 服务器 headful 抓取配置(代理/超时) |
| `configs/profiles/8g_headful.yaml` | Pi 服务器 headful 抓取配置(代理/超时) | | `configs/profiles/8g_headful.yaml` | Pi 服务器 headful 抓取配置(代理/超时) |
| `.env` | 密钥 / 服务地址(不入 Git) | | `.env` | 密钥 / 服务地址(不入 Git) |
| `prompts/` | LLM Prompt 模板(翻译/日报/搜索 Agent) | | `prompts/` | LLM Prompt 模板(翻译/日报/搜索 Agent) |
### AI 模型按场景配置(llm_scenes)
大模型按场景独立配置,见 `configs/system.yaml` 的 `llm_scenes` 段:
| 场景 | 用途 | 模型(当前) | 参数 |
|------|------|-------------|------|
| `translation` | M4 全文英译中 + 投资事件抽取 | deepseek-v4-flash | temperature=0.1, max_tokens=8192 |
| `daily_report` | M7 日报 AI 摘要(分批生成) | deepseek-v4-flash | temperature=0.3, max_tokens=1500 |
场景未声明的字段回退 `llm` 默认段;Embedding 为单一场景(`en_finance_news` 库入库/检索向量必须同模型,不支持拆分)。
### 去重多来源(M3)
去重时跨源重复的新闻,会把所有来源记录到保留的唯一篇 `source_ids` 字段(首个来源为 `source_id`),经翻译透传后:
- `news_event.source`:拼接展示名(如 "Barron's, CNBC, Reuters",最多前 3 个,向后兼容)
- `news_event.sources`(TEXT,JSON 数组):**全部来源展示名,主源居首**,如 `["Barron's","CNBC","Reuters","Financial Times"]`,前端直接渲染完整来源列表
表结构变更(`news_event` 新增 `sources` 列)由 `report_db.init_schema()` 幂等迁移(`ALTER TABLE ... ADD COLUMN IF NOT EXISTS`,MariaDB 10.0.2+)。
### 中断恢复(--resume)
全流程脚本支持断点续跑:
- `scripts/domestic_full.sh --resume` / `scripts/pipeline.sh --resume`:跳过当天已完成步骤(步骤状态存 `data/run_state/{YYYYMMDD}.state`,按天换新,跨天自动失效)
- 步骤粒度:`M1_crawl` / `M2_extract` / `M3_dedup` / `M4_translate` / `M5_embed` / `M6_index` / `report`;某步失败不标记,`--resume` 从失败处重试
- **文件级增量(各步骤内自动跳过已处理文件)**:M2 按 `data/processed/.../{url_hash}.json` 存在性、M3 按指纹库判重、M4 按 `data/events/.../{url_hash}.json` 存在性、M5 按 `data/embeddings/.../index.json`、M6 按 Qdrant upsert 幂等(点 id=url_hash)、日报按 MySQL 唯一键覆盖
### 终端进度与日志
执行全流程时,终端**实时显示完整进度**,且**显性标注当前阶段与 AI 模型**:
- 每阶段标题:`━━━ M2 正文提取 ━━━`、`━━━ M4 翻译+事件抽取(AI 大模型: deepseek / deepseek-v4-flash)━━━`(AI 阶段自动从 `system.yaml` 读取供应商/模型)
- 每步进度:`▶ 步骤 开始` → 步骤内逐源/逐篇输出 → `✔ 完成(耗时 Ns)`;失败显示 `✗` 与退出码
- AI 调用点在 Python 层同样显性打印(`AI 大模型(场景 translation/daily_report): provider=... model=...`、`初始化 Embedding 客户端: provider=dashscope model=...`)
- 日志:`logs/domestic_full_{ts}.log`(全流程)、`logs/pipeline_{ts}.log`(管道),已入 `.gitignore`
### 日报入库(M9) ### 日报入库(M9)
日报内容结构化写入与 [news 项目](https://github.com/) 共用的 MySQL `myquant` 库(表 `news_report` / `news_event`,`report_type="intl"`,同一天重复生成幂等覆盖)。表结构与数据契约见 news 项目 `docs/db_schema.md`。 日报内容结构化写入与 [news 项目](https://github.com/) 共用的 MySQL `myquant` 库(表 `news_report` / `news_event`,`report_type="intl"`,同一天重复生成幂等覆盖)。表结构与数据契约见 news 项目 `docs/db_schema.md`。
+43 -4
View File
@@ -40,23 +40,62 @@ dedup:
simhash_window_days: 30 simhash_window_days: 30
min_content_length: 100 min_content_length: 100
# ── LLM 翻译+事件抽取 ──────────────────────────────── # ── LLM 默认配置(所有 LLM 场景的兜底)────────────────
# 按场景独立配置见下方 llm_scenes 段;场景未声明的字段回退到本段。
llm: llm:
provider: "deepseek" provider: "deepseek"
deepseek_model: "deepseek-chat" deepseek_model: "deepseek-v4-flash"
qwen_model: "qwen-plus" qwen_model: "qwen-plus"
timeout_sec: 60 timeout_sec: 60
max_retries: 3 max_attempts: 3 # 单篇总尝试次数(含首次),失败后指数退避重试
max_tokens: 8192 max_tokens: 8192
temperature: 0.1 temperature: 0.1
concurrency: 3 concurrency: 3
# ── Embedding 向量化 ──────────────────────────────── # ── LLM 场景配置(按场景独立指定大模型类型与参数)──────
# 每个场景可覆盖 provider / model / temperature / max_tokens / max_attempts / timeout_sec;
# 场景内统一用 "model" 键指定模型(优先于 llm 段的 deepseek_model / qwen_model)。
llm_scenes:
translation:
# M4:全文英译中 + 投资事件抽取(单次 LLM 调用合并输出)
provider: "deepseek"
model: "deepseek-v4-flash"
temperature: 0.1
max_tokens: 8192
description: |
用途: M4 对去重后的英文正文做全文英译中,并抽取投资事件
(事件类型/美股代码/情绪/重要度/摘要,单次调用合并输出)
使用方法: llm/pipeline.py 调用 load_llm_config(scene="translation"),
配合 concurrency=3 逐篇并发;单篇失败重试 max_attempts 次后跳过,
未翻译篇由增量机制下次补齐
模型要求: 中英财经翻译准确、术语一致;严格按 Prompt 输出 JSON 结构;
单篇平均 ≤3 秒;max_tokens 需容纳长文(建议 ≥8192)
daily_report:
# M7:每日 AI 摘要日报(五段式,高重要度事件分批生成)
provider: "deepseek"
model: "deepseek-v4-flash"
temperature: 0.3
max_tokens: 1500
description: |
用途: M7 日报 AI 摘要(五段式),高重要度事件分批(≤10 条/批)生成
各批摘要后合并为完整日报摘要
使用方法: scheduler/reporter.py::_call_llm_simple 使用
load_llm_config(scene="daily_report");分批/合并失败均有回退,
全部失败走规则兜底(直接列 Top 事件),日报仍正常入库
模型要求: 中文财经总结能力强;单批 600-1500 字摘要质量稳定;
支持高频短调用(crontab 07/12/18 每天 3 次 × 每份 3 批)
# ── Embedding 向量化(单一场景,不按场景拆分)─────────
# 说明: Qdrant collection en_finance_news 的入库向量与检索查询向量必须由
# 同一模型生成(跨模型向量无法比较),因此 embedding 不支持按场景独立配置。
# 调用点: embedding/pipeline.py(入库)、vectorstore/pipeline.py(检索)、
# mcp_server/server.py(MCP 搜索查询向量化)——三处共用本配置。
embedding: embedding:
provider: "dashscope" provider: "dashscope"
dashscope_model: "text-embedding-v3" dashscope_model: "text-embedding-v3"
dimension: 1024 dimension: 1024
batch_size: 10 batch_size: 10
max_attempts: 3 # 单批总尝试次数(含首次),失败后指数退避重试
timeout_sec: 30 timeout_sec: 30
# ── Qdrant ────────────────────────────────────────── # ── Qdrant ──────────────────────────────────────────
+113 -1
View File
@@ -1,6 +1,118 @@
# continuation.md — English Financial News 项目状态 # continuation.md — English Financial News 项目状态
> 最后更新:2026-08-04 > 最后更新:2026-08-12
---
## 2026-08-12 会话成果(五):日报多来源入库(news_event.sources 列)
**背景:** 前端需要展示一条新闻的多个来源;`source` 拼接字段(≤3 个)不够。
**改动(用户已授权改表,列已由用户添加):**
| 文件 | 改动内容 |
|------|---------|
| `report_db/models.py` | `EventRow` 新增 `sources: list[str]`(全部来源展示名,主源居首) |
| `report_db/schema.py` | DDL 加 `sources TEXT NULL` 列;`init_schema` 增加幂等 `ALTER TABLE ... ADD COLUMN IF NOT EXISTS sources`(旧表迁移) |
| `report_db/db.py` | `save_report` INSERT 写入 sources(JSON 数组) |
| `scheduler/reporter.py` | 新增 `_article_sources_full()`(完整来源列表,主源居首);`_article_source_label` 重构为基于它(source 保持 ≤3 拼接兼容);`_build_report_data` 事件写入 `sources` |
| `tests/test_report_db.py` | 多源完整列表断言(2 源 / 4 源截断对比 / 单源回退) |
**DB 迁移:** `news_event.sources` TEXT 列(用户已加);`init_schema()` 幂等补列(MariaDB `ADD COLUMN IF NOT EXISTS`),pi5 实测幂等 OK。
**部署中踩坑记录:** `rsync -a report_db/ scheduler/reporter.py pi5:/home/pi/intlnews/` 的混合参数把文件散落到目标根目录(report_db/ 带斜杠 = 内容进根、reporter.py 也进了根),导致 pi5 上模块未更新、sources 写入 None。**教训:rsync 多源参数混用目录与文件时目标路径易错,部署后必须 grep 关键标识(如 source_list)验证。** 已清理根目录残留并正确同步。
**验证:** 单元测试 3 用例(多源完整列表);pi5 实盘:`init_schema` 幂等、日报 report_id=224 sources 列全部写入(`["InvestingLive"]` 单源数组)。当前真实数据无多源新闻入日报(跨源合并发生在历史日期文件,未重新翻译),多源数组待后续新数据自然出现(链路已通,前端可直接读 `sources` JSON 数组)。
---
## 2026-08-12 会话成果(四):脚本阶段/模型显性输出
**背景:** 执行全流程时终端未显性标注"当前阶段",AI 调用的供应商/模型只在客户端初始化时输出一次。
**改动:**
| 文件 | 改动内容 |
|------|---------|
| `scripts/pipeline.sh` | 每阶段 banner(`━━━ M4 翻译+事件抽取(AI 大模型: deepseek / deepseek-v4-flash)━━━`);`ai_model_info()` 从 system.yaml 读取场景模型(translation/daily_report/embedding) |
| `scripts/domestic_full.sh` | M1 / 管道阶段标题统一 banner 风格 |
| `llm/pipeline.py` | 创建客户端后显性输出 `AI 大模型(场景 translation): provider/model` |
| `embedding/pipeline.py` | 向量化日志补充 provider |
| `scheduler/reporter.py` | daily_report 场景客户端初始化时输出 provider/model |
**测试:** 本地模拟 banner 全部正确(M4/M5/日报 3 处 AI 标注);pi5 实测客户端初始化日志含 provider/model(deepseek/deepseek-v4-flash、dashscope/text-embedding-v3);全量 184 passed。
---
## 2026-08-12 会话成果(三):全流程终端进度输出
**背景:** `domestic_full.sh` / `pipeline.sh` 原用 `tail -5/-10` 截断输出,终端看不到中间进度。
**改动:**
| 文件 | 改动内容 |
|------|---------|
| `scripts/_step_state.sh` | `step_run` 增加耗时统计与进度日志(`▶ 开始` / `✔ 完成(耗时 Ns)` / `✗ 失败(exit=N,耗时)`);支持 `STEP_LOG` 环境变量——步骤全部输出实时显示终端并同时写入日志(用 `PIPESTATUS[0]` 取原命令退出码,规避 pipefail/tee 吞退出码) |
| `scripts/pipeline.sh` | 步骤输出透传终端(去掉 `tail -10`);日志写 `logs/pipeline_{ts}.log` |
| `scripts/domestic_full.sh` | M1 与管道输出透传终端(去掉 `tail -5/-10`);日志写 `logs/domestic_full_{ts}.log` |
**测试:** 无 pipefail / pipefail+条件 / pipefail+非条件 三环境下退出码捕获均正确(失败不 mark);pi5 实机 `--resume` 跑通(M3 31s、M6 20s 均显示耗时,日志文件生成)。
**注意:** 所有含 `${var}` 后跟全角标点的输出已用花括号包裹(bash 会把全角字符字节并入变量名,导致 unbound variable)。
---
## 2026-08-12 会话成果(二):全流程 --resume 断点续跑
**背景:** `scripts/domestic_full.sh`(M1→M6→日报)中断(ssh 断线/断电)后只能从头重跑。
**改动:**
| 文件 | 改动内容 |
|------|---------|
| `scripts/_step_state.sh`(新增) | 步骤状态库:`step_run` / `step_should_run` / `step_mark` / `step_done`;状态文件 `data/run_state/{YYYYMMDD}.state`(按天换新);`step_run` 显式检查退出码(规避 bash set -e 条件上下文陷阱) |
| `scripts/pipeline.sh` | 支持 `--resume`;M2→M6→report 每步 `step_run` 包裹 |
| `scripts/domestic_full.sh` | 支持 `--resume`;M1 用 `step_should_run`/`step_mark`(部分源失败仍标记,原语义);透传 `--resume` 给 pipeline.sh |
| `.gitignore` | 忽略 `data/run_state/` |
**文件级增量确认(各步骤已内置,无需新写):** M2 processed JSON 存在跳过、M3 指纹库判重、M4 events JSON 存在跳过、M5 embedding index、M6 Qdrant upsert 幂等(id=url_hash)、日报 MySQL 唯一键覆盖。
**测试:** 本地模拟 3 场景全过(失败步骤不 mark / set -e 退出不 mark / resume 跳过已完成);`--badarg` exit=1;pi5 同步后验证 step_run 机制与参数校验正常。
**注意:** crontab 07/12/18 跑的是不带 --resume 的全量(步骤级不跳过,靠文件级增量),--resume 供手动中断恢复。
---
## 2026-08-12 会话成果
### M9.2:AI 模型按场景配置 + 去重多来源
**背景:** ① 本项目所有 AI 大模型调用点(M4 翻译+事件抽取、M7 日报 AI 摘要、M5 向量化)此前共用同一套 `llm`/`embedding` 配置;② M3 去重跨源重复时丢弃重复篇来源信息。
**改动:**
| 文件 | 改动内容 |
|------|---------|
| `configs/system.yaml` | 新增 `llm_scenes` 段(translation / daily_report,含用途/使用方法/模型要求说明,字段回退 `llm` 默认段);`embedding` 段注释说明单一场景原因 |
| `llm/client.py` | `_load_system_config(scene)` 场景合并;`load_llm_config(..., scene)` 支持按场景覆盖 provider/model/参数;场景统一用 `model` 键 |
| `llm/pipeline.py` | `load_llm_config(scene="translation")` |
| `scheduler/reporter.py` | `_call_llm_simple` 用 `scene="daily_report"`;`temperature` 硬编码 0.3 → `config.temperature`(技术债清除);新增 `_article_source_label()` 多来源拼接展示 |
| `extractor/models.py` / `llm/models.py` | `ProcessedArticle` / `EnTranslatedArticle` 新增 `source_ids` 字段 |
| `dedup/pipeline.py` | unique 初始化 `source_ids=[source_id]`;dup 时 `_merge_duplicate_source()` 跨日期目录合并来源进唯一篇 |
| `llm/extractor.py` | 同步/异步构造 `EnTranslatedArticle` 透传 `source_ids` |
| `tests/` | `test_llm.py` 场景配置 4 用例;`test_dedup.py::TestMergeSources` 3 用例;`test_report_db.py` 多来源拼接 2 用例 |
**测试结果:** 本地全量 184 passed / 2 failed(原有 test_crawler 路径问题);ruff 无新增。
**部署验证(pi5 实盘):**
- scene 配置生效:translation=(v4-flash, 0.1)、daily_report=(v4-flash, 0.3, 1500)
- 去重合并实证:历史唯一篇 `049da7e0c87160dc`(barrons 主源)被 9 个跨源重复篇合并 → `source_ids` 10 个来源
- 日报仍正常入库:report_id=224(2026-08-12, 15 事件);当天 DeepSeek 摘要 3 次返回空 → 规则兜底,日报未中断(失败处理按设计工作)
- 多来源展示:`_article_source_label` 实测 "Barron's, CNBC, Reuters";单源回退正常
**已知说明:**
- 历史 3368 个 deduped 旧文件无 `source_ids` 字段(不回填,向前生效);events 文件由下次 crontab 07:00 用新代码自然带出透传
- 179 篇全重复(增量正常);138 个跨源重复的 merge 目标多为历史日期目录文件(指纹库 ±30 天窗口所致),非当天 uniques
--- ---
+37
View File
@@ -103,6 +103,8 @@ def dedup_source(
if result.is_duplicate: if result.is_duplicate:
dup_count += 1 dup_count += 1
# 跨源重复:把来源合并进已保留的唯一篇(记录多个来源)
_merge_duplicate_source(article, result)
logger.debug("[%s] 🔁 %s → L%d: %s", logger.debug("[%s] 🔁 %s → L%d: %s",
source_id, source_id,
article.title[:40], article.title[:40],
@@ -110,6 +112,9 @@ def dedup_source(
result.short_summary()) result.short_summary())
else: else:
unique_count += 1 unique_count += 1
# 初始化来源列表(首个来源 = 本篇文章来源)
if not article.source_ids:
article.source_ids = [article.source_id]
# 写入唯一条目 # 写入唯一条目
out_file = out_dir / f"{article.url_hash}.json" out_file = out_dir / f"{article.url_hash}.json"
out_file.write_text( out_file.write_text(
@@ -130,6 +135,38 @@ def dedup_source(
} }
def _merge_duplicate_source(article: ProcessedArticle, result: DedupResult) -> None:
"""重复篇:把来源 ID 追加进已保留的唯一篇 JSON(最终显示的新闻记录多个来源)。
唯一篇文件按 url_hash 定位(跨日期目录搜索,因指纹窗口为 ±30 天);
文件不存在(超窗口被清理)时仅记录日志,不阻塞去重流程。
"""
if not result.matched_url_hash:
return
candidates = sorted(Path("data/deduped").glob(f"*/uniques/{result.matched_url_hash}.json"))
if not candidates:
logger.warning(
"重复篇唯一文件不存在(可能已超窗口): %s(重复来源 %s 未合并)",
result.matched_url_hash, article.source_id,
)
return
target = candidates[0]
try:
data = json.loads(target.read_text(encoding="utf-8"))
# 旧格式文件可能无 source_ids:以主来源 source_id 兜底
merged = list(dict.fromkeys(
[*(data.get("source_ids") or [data.get("source_id")]), article.source_id]
))
data["source_ids"] = merged
target.write_text(
json.dumps(data, indent=2, ensure_ascii=False), encoding="utf-8"
)
logger.debug("来源合并: %s → %s (sources=%s)",
article.source_id, result.matched_url_hash, merged)
except Exception as e:
logger.exception("来源合并失败 %s: %s", target, e)
def _layer_num(result: DedupResult) -> int: def _layer_num(result: DedupResult) -> int:
"""DedupResult → 命中层编号。""" """DedupResult → 命中层编号。"""
if result.matched_layer is None: if result.matched_layer is None:
+2
View File
@@ -92,6 +92,7 @@ def load_embedding_config(
dimension = int(sys_cfg.get("dimension", DASHSCOPE_DEFAULT_DIM)) dimension = int(sys_cfg.get("dimension", DASHSCOPE_DEFAULT_DIM))
batch_size = min(int(sys_cfg.get("batch_size", DASHSCOPE_BATCH_LIMIT)), DASHSCOPE_BATCH_LIMIT) batch_size = min(int(sys_cfg.get("batch_size", DASHSCOPE_BATCH_LIMIT)), DASHSCOPE_BATCH_LIMIT)
timeout = float(sys_cfg.get("timeout_sec", 30.0)) timeout = float(sys_cfg.get("timeout_sec", 30.0))
max_attempts = int(sys_cfg.get("max_attempts", DEFAULT_MAX_ATTEMPTS))
if not api_key: if not api_key:
raise EmbeddingError("DASHSCOPE_API_KEY 未配置,请检查 .env") raise EmbeddingError("DASHSCOPE_API_KEY 未配置,请检查 .env")
@@ -104,6 +105,7 @@ def load_embedding_config(
dimension=dimension, dimension=dimension,
batch_size=batch_size, batch_size=batch_size,
timeout_sec=timeout, timeout_sec=timeout,
max_attempts=max_attempts,
) )
+2 -2
View File
@@ -100,8 +100,8 @@ def embed_all_events(
# 批量嵌入(按 batch_size 分块,每批输出进度) # 批量嵌入(按 batch_size 分块,每批输出进度)
batch_size = config.batch_size batch_size = config.batch_size
total = len(articles) total = len(articles)
logger.info("开始向量化 %d 篇文章(batch_size=%d, model=%s)", logger.info("开始向量化 %d 篇文章(batch_size=%d, provider=%s, model=%s)",
total, batch_size, config.model) total, batch_size, config.provider, config.model)
for batch_start in range(0, total, batch_size): for batch_start in range(0, total, batch_size):
batch_end = min(batch_start + batch_size, total) batch_end = min(batch_start + batch_size, total)
+6 -1
View File
@@ -1,6 +1,6 @@
"""正文提取数据模型""" """正文提取数据模型"""
from pydantic import BaseModel from pydantic import BaseModel, Field
class ProcessedArticle(BaseModel): class ProcessedArticle(BaseModel):
@@ -20,3 +20,8 @@ class ProcessedArticle(BaseModel):
status: str = "success" # success | no_content | failed status: str = "success" # success | no_content | failed
extractor: str = "trafilatura" # trafilatura | crawl4ai_md | none extractor: str = "trafilatura" # trafilatura | crawl4ai_md | none
error: str = "" error: str = ""
source_ids: list[str] = Field(
default_factory=list,
description="去重合并后的所有来源 ID(M3 去重时跨源命中重复会追加;"
"首个来源始终为 source_id)",
)
+28 -8
View File
@@ -25,15 +25,30 @@ _QWEN_DEFAULT_BASE = "https://dashscope.aliyuncs.com/compatible-mode/v1"
_DEEPSEEK_DEFAULT_MODEL = "deepseek-chat" _DEEPSEEK_DEFAULT_MODEL = "deepseek-chat"
_QWEN_DEFAULT_MODEL = "qwen-plus" _QWEN_DEFAULT_MODEL = "qwen-plus"
# 默认单次调用总尝试次数(含首次;system.yaml llm.max_attempts 未配置时兜底)
_DEFAULT_MAX_ATTEMPTS = 3
def _load_system_config() -> dict:
"""加载 configs/system.yaml 中 llm 段配置。""" def _load_system_config(scene: str | None = None) -> dict:
"""加载 llm 配置;scene 指定时与 llm_scenes.{scene} 合并(场景覆盖默认段)。
Args:
scene: 场景名(translation / daily_report)。未配置该场景时回退 llm 段。
"""
config_path = Path("configs/system.yaml") config_path = Path("configs/system.yaml")
if config_path.exists(): if config_path.exists():
try: try:
with open(config_path, encoding="utf-8") as f: with open(config_path, encoding="utf-8") as f:
raw = yaml.safe_load(f) raw = yaml.safe_load(f)
return raw.get("llm", {}) base = raw.get("llm", {})
if scene:
scene_cfg = (raw.get("llm_scenes", {}) or {}).get(scene, {})
if not scene_cfg:
logger.warning(
"system.yaml 中不存在 llm_scenes.%s,使用 llm 默认段", scene
)
return {**base, **scene_cfg}
return base
except Exception: except Exception:
logger.warning("加载 llm 配置失败,使用空配置") logger.warning("加载 llm 配置失败,使用空配置")
return {} return {}
@@ -50,6 +65,7 @@ class LLMConfig:
timeout_sec: float = 60.0 timeout_sec: float = 60.0
temperature: float = 0.1 temperature: float = 0.1
max_tokens: int = 8192 max_tokens: int = 8192
max_attempts: int = _DEFAULT_MAX_ATTEMPTS # 单次调用总尝试次数(含首次)
def __post_init__(self) -> None: def __post_init__(self) -> None:
if not self.api_key: if not self.api_key:
@@ -60,26 +76,28 @@ def load_llm_config(
provider: str | None = None, provider: str | None = None,
*, *,
model: str | None = None, model: str | None = None,
scene: str | None = None,
) -> LLMConfig: ) -> LLMConfig:
"""根据配置文件构造 LLMConfig。 """根据配置文件构造 LLMConfig。
provider 为 None 时读 system.yaml llm.provider,默认 deepseek。 provider 为 None 时读 system.yaml(或场景段)llm.provider,默认 deepseek。
model 为 None 时读 system.yaml 中对应 provider 的 model。 model 为 None 时优先读场景段的 model 键,其次 system.yaml 对应 provider 的 model。
scene 指定时,llm_scenes.{scene} 覆盖 llm 默认段的各参数(按场景独立配置模型)。
Raises: Raises:
ValueError: API key 未配置 ValueError: API key 未配置
""" """
config = _load_system_config() config = _load_system_config(scene=scene)
p = (provider or config.get("provider", "deepseek")).lower() p = (provider or config.get("provider", "deepseek")).lower()
if p == "deepseek": if p == "deepseek":
api_key = os.environ.get("DEEPSEEK_API_KEY", "") api_key = os.environ.get("DEEPSEEK_API_KEY", "")
base = os.environ.get("DEEPSEEK_BASE_URL", _DEEPSEEK_DEFAULT_BASE) base = os.environ.get("DEEPSEEK_BASE_URL", _DEEPSEEK_DEFAULT_BASE)
m = model or config.get("deepseek_model", _DEEPSEEK_DEFAULT_MODEL) m = model or config.get("model") or config.get("deepseek_model", _DEEPSEEK_DEFAULT_MODEL)
elif p in ("qwen", "dashscope"): elif p in ("qwen", "dashscope"):
api_key = os.environ.get("QWEN_API_KEY") or os.environ.get("DASHSCOPE_API_KEY") or "" api_key = os.environ.get("QWEN_API_KEY") or os.environ.get("DASHSCOPE_API_KEY") or ""
base = os.environ.get("QWEN_BASE_URL", _QWEN_DEFAULT_BASE) base = os.environ.get("QWEN_BASE_URL", _QWEN_DEFAULT_BASE)
m = model or config.get("qwen_model", _QWEN_DEFAULT_MODEL) m = model or config.get("model") or config.get("qwen_model", _QWEN_DEFAULT_MODEL)
p = "qwen" p = "qwen"
else: else:
raise ValueError(f"未知 LLM provider: {p!r},仅支持 deepseek / qwen") raise ValueError(f"未知 LLM provider: {p!r},仅支持 deepseek / qwen")
@@ -93,6 +111,7 @@ def load_llm_config(
timeout = float(config.get("timeout_sec", 60.0)) timeout = float(config.get("timeout_sec", 60.0))
temperature = float(config.get("temperature", 0.1)) temperature = float(config.get("temperature", 0.1))
max_tokens = int(config.get("max_tokens", 8192)) max_tokens = int(config.get("max_tokens", 8192))
max_attempts = int(config.get("max_attempts", _DEFAULT_MAX_ATTEMPTS))
return LLMConfig( return LLMConfig(
provider=p, provider=p,
@@ -102,6 +121,7 @@ def load_llm_config(
timeout_sec=timeout, timeout_sec=timeout,
temperature=temperature, temperature=temperature,
max_tokens=max_tokens, max_tokens=max_tokens,
max_attempts=max_attempts,
) )
+10 -3
View File
@@ -226,7 +226,7 @@ def translate_and_extract(
article: ProcessedArticle, article: ProcessedArticle,
*, *,
template: PromptTemplate | None = None, template: PromptTemplate | None = None,
max_attempts: int = DEFAULT_MAX_ATTEMPTS, max_attempts: int | None = None,
) -> EnTranslatedArticle: ) -> EnTranslatedArticle:
"""同步翻译 + 事件抽取(单篇文章,带重试)。 """同步翻译 + 事件抽取(单篇文章,带重试)。
@@ -235,7 +235,8 @@ def translate_and_extract(
config: LLM 配置 config: LLM 配置
article: 待处理的英文新闻 article: 待处理的英文新闻
template: Prompt 模板,默认加载 prompts/translation_and_extraction.md template: Prompt 模板,默认加载 prompts/translation_and_extraction.md
max_attempts: 最大重试次数 max_attempts: 最大尝试次数(含首次);None 时取 config.max_attempts
(来自 system.yaml llm.max_attempts)
Returns: Returns:
EnTranslatedArticle 含双语内容 + 事件 EnTranslatedArticle 含双语内容 + 事件
@@ -243,6 +244,8 @@ def translate_and_extract(
Raises: Raises:
LLMCallError: 所有重试均失败 LLMCallError: 所有重试均失败
""" """
if max_attempts is None:
max_attempts = config.max_attempts
tpl = template or PromptTemplate() tpl = template or PromptTemplate()
system_prompt, user_prompt = tpl.render(article) system_prompt, user_prompt = tpl.render(article)
@@ -261,6 +264,7 @@ def translate_and_extract(
source_name=article.source_name, source_name=article.source_name,
url=article.url, url=article.url,
url_hash=article.url_hash, url_hash=article.url_hash,
source_ids=list(article.source_ids),
title=article.title, title=article.title,
title_zh=output.title_zh, title_zh=output.title_zh,
content_en=article.content, content_en=article.content,
@@ -305,10 +309,12 @@ async def translate_and_extract_async(
article: ProcessedArticle, article: ProcessedArticle,
*, *,
template: PromptTemplate | None = None, template: PromptTemplate | None = None,
max_attempts: int = DEFAULT_MAX_ATTEMPTS, max_attempts: int | None = None,
semaphore: asyncio.Semaphore | None = None, semaphore: asyncio.Semaphore | None = None,
) -> EnTranslatedArticle: ) -> EnTranslatedArticle:
"""异步翻译 + 事件抽取(批处理用),与同步版逻辑等价。""" """异步翻译 + 事件抽取(批处理用),与同步版逻辑等价。"""
if max_attempts is None:
max_attempts = config.max_attempts
tpl = template or PromptTemplate() tpl = template or PromptTemplate()
system_prompt, user_prompt = tpl.render(article) system_prompt, user_prompt = tpl.render(article)
@@ -327,6 +333,7 @@ async def translate_and_extract_async(
source_name=article.source_name, source_name=article.source_name,
url=article.url, url=article.url,
url_hash=article.url_hash, url_hash=article.url_hash,
source_ids=list(article.source_ids),
title=article.title, title=article.title,
title_zh=output.title_zh, title_zh=output.title_zh,
content_en=article.content, content_en=article.content,
+4
View File
@@ -103,6 +103,10 @@ class EnTranslatedArticle(BaseModel):
source_name: str source_name: str
url: str url: str
url_hash: str url_hash: str
source_ids: list[str] = Field(
default_factory=list,
description="去重合并后的所有来源 ID(透传自 ProcessedArticle.source_ids)",
)
# ── 双语内容 ── # ── 双语内容 ──
title: str = "" # 英文原标题 title: str = "" # 英文原标题
+6 -2
View File
@@ -119,9 +119,13 @@ def translate_all_deduped(
logger.warning("去重目录无文章: data/deduped/%s/uniques/", date_str) logger.warning("去重目录无文章: data/deduped/%s/uniques/", date_str)
return {"date": date_str, "total": 0, "success": 0, "failed": 0, "elapsed_sec": 0} return {"date": date_str, "total": 0, "success": 0, "failed": 0, "elapsed_sec": 0}
# 初始化 LLM 客户端 # 初始化 LLM 客户端(translation 场景配置见 system.yaml llm_scenes.translation)
config = load_llm_config(provider=provider, model=model) config = load_llm_config(provider=provider, model=model, scene="translation")
client = make_sync_client(config) client = make_sync_client(config)
logger.info(
"AI 大模型(场景 translation): provider=%s model=%s",
config.provider, config.model,
)
template = PromptTemplate() template = PromptTemplate()
# 输出目录 # 输出目录
+4 -2
View File
@@ -113,12 +113,13 @@ def save_report(conn: pymysql.Connection, report: ReportData) -> int:
cur.execute("DELETE FROM news_event WHERE report_id=%s", (report_id,)) cur.execute("DELETE FROM news_event WHERE report_id=%s", (report_id,))
for ev in report.events: for ev in report.events:
sources_json = json.dumps(ev.sources, ensure_ascii=False) if ev.sources else None
cur.execute( cur.execute(
""" """
INSERT INTO news_event INSERT INTO news_event
(report_id, section, rank, importance, event_type, title, (report_id, section, rank, importance, event_type, title,
summary, sentiment, source, url) summary, sentiment, source, sources, url)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
""", """,
( (
report_id, report_id,
@@ -130,6 +131,7 @@ def save_report(conn: pymysql.Connection, report: ReportData) -> int:
ev.summary, ev.summary,
ev.sentiment, ev.sentiment,
ev.source, ev.source,
sources_json,
ev.url, ev.url,
), ),
) )
+4
View File
@@ -22,6 +22,10 @@ class EventRow(BaseModel):
summary: str | None = None summary: str | None = None
sentiment: str | None = None # positive | negative | neutral | '' sentiment: str | None = None # positive | negative | neutral | ''
source: str | None = None source: str | None = None
sources: list[str] = Field(
default_factory=list,
description="全部来源展示名(JSON 数组,主源居首;多源新闻记录完整列表)",
)
url: str | None = None url: str | None = None
+8 -1
View File
@@ -34,11 +34,18 @@ DDL_STATEMENTS: list[str] = [
title VARCHAR(512) NOT NULL COMMENT '标题', title VARCHAR(512) NOT NULL COMMENT '标题',
summary TEXT NULL COMMENT '摘要/正文', summary TEXT NULL COMMENT '摘要/正文',
sentiment VARCHAR(8) NULL COMMENT 'positive/negative/neutral', sentiment VARCHAR(8) NULL COMMENT 'positive/negative/neutral',
source VARCHAR(64) NULL COMMENT '来源', source VARCHAR(64) NULL COMMENT '来源(多源时拼接展示名,最多前 3 个)',
sources TEXT NULL COMMENT '全部来源(JSON 数组,主源居首;多源新闻完整列表)',
url VARCHAR(512) NULL COMMENT '原文链接', url VARCHAR(512) NULL COMMENT '原文链接',
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
KEY idx_report_section (report_id, section), KEY idx_report_section (report_id, section),
KEY idx_title (title(255)) KEY idx_title (title(255))
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='日报事件明细' ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='日报事件明细'
""", """,
# 旧表补列(幂等;MariaDB 10.0.2+ 支持 ADD COLUMN IF NOT EXISTS)
"""
ALTER TABLE news_event
ADD COLUMN IF NOT EXISTS sources TEXT NULL
COMMENT '全部来源(JSON 数组,主源居首;多源新闻完整列表)' AFTER source
""",
] ]
+45 -7
View File
@@ -96,6 +96,33 @@ def _url_source_label(url: str, source_id: str = "") -> str:
return domain_map.get(domain, domain) return domain_map.get(domain, domain)
def _article_sources_full(article: dict) -> list[str]:
"""全部来源展示名列表(主源居首、去重)。
多来源(ProcessedArticle.source_ids)时返回完整列表供 news_event.sources
JSON 数组使用;无合并信息时回退单源逻辑。
"""
src_ids = list(dict.fromkeys(
s for s in (article.get("source_ids") or []) if s and s != "?"
))
if len(src_ids) > 1:
return list(dict.fromkeys((_source_name(s) or s) for s in src_ids))
src = _url_source_label(article.get("url"), article.get("source_id", ""))
return [src] if src not in (None, "", "?") else []
def _article_source_label(article: dict) -> str | None:
"""文章来源展示(向后兼容 news_event.source):多来源时拼接前 3 个展示名。
多来源(source_ids 长度 > 1)时展示如 "Reuters, CNBC",最多取前 3 个来源,
截断至 64 字符(source 字段为 VARCHAR(64));完整列表见 _article_sources_full。
"""
names = _article_sources_full(article)
if len(names) > 1:
return ", ".join(names[:3])[:64]
return names[0] if names else None
# 日报覆盖时间窗口(小时) # 日报覆盖时间窗口(小时)
_REPORT_WINDOW_HOURS = 25 _REPORT_WINDOW_HOURS = 25
@@ -283,7 +310,7 @@ def _call_llm_simple(
system_prompt: str, system_prompt: str,
user_prompt: str, user_prompt: str,
max_tokens: int = 600, max_tokens: int = 600,
max_retries: int = 2, max_retries: int | None = None,
) -> str: ) -> str:
"""封装 LLM 调用,带重试和客户端复用。 """封装 LLM 调用,带重试和客户端复用。
@@ -291,23 +318,30 @@ def _call_llm_simple(
system_prompt: system role 内容 system_prompt: system role 内容
user_prompt: user role 内容 user_prompt: user role 内容
max_tokens: 最大输出 token max_tokens: 最大输出 token
max_retries: 最大重试次数(不含首次调用) max_retries: 重试次数(不含首次);None 时取
config.max_attempts - 1(system.yaml llm.max_attempts)
Returns: Returns:
LLM 输出文本;所有重试均失败返回空字符串 LLM 输出文本;所有重试均失败返回空字符串
""" """
import time as _time import time as _time
# 复用客户端(同 provider/model 只创建一次) # 复用客户端(同 provider/model 只创建一次;日报摘要场景见 system.yaml llm_scenes.daily_report)
cache_key = "default" cache_key = "default"
if cache_key not in _llm_client_cache: if cache_key not in _llm_client_cache:
from llm.client import load_llm_config, make_sync_client from llm.client import load_llm_config, make_sync_client
_llm_client_cache["config"] = load_llm_config() _llm_client_cache["config"] = load_llm_config(scene="daily_report")
_llm_client_cache[cache_key] = make_sync_client(_llm_client_cache["config"]) _llm_client_cache[cache_key] = make_sync_client(_llm_client_cache["config"])
logger.info("AI 大模型(场景 daily_report): provider=%s model=%s",
_llm_client_cache["config"].provider,
_llm_client_cache["config"].model)
config = _llm_client_cache["config"] config = _llm_client_cache["config"]
client = _llm_client_cache[cache_key] client = _llm_client_cache[cache_key]
if max_retries is None:
max_retries = max(config.max_attempts - 1, 0)
last_err: str = "" last_err: str = ""
for attempt in range(1, max_retries + 2): # 首次 + max_retries 次重试 for attempt in range(1, max_retries + 2): # 首次 + max_retries 次重试
try: try:
@@ -317,7 +351,7 @@ def _call_llm_simple(
{"role": "system", "content": system_prompt}, {"role": "system", "content": system_prompt},
{"role": "user", "content": user_prompt}, {"role": "user", "content": user_prompt},
], ],
temperature=0.3, temperature=config.temperature,
max_tokens=max_tokens, max_tokens=max_tokens,
) )
content = (resp.choices[0].message.content or "").strip() content = (resp.choices[0].message.content or "").strip()
@@ -577,7 +611,10 @@ def _build_report_data(
article = ev.get("article", {}) article = ev.get("article", {})
title = (article.get("title_zh") or article.get("title") or "").strip()[:512] title = (article.get("title_zh") or article.get("title") or "").strip()[:512]
url = article.get("url") or None url = article.get("url") or None
src = _url_source_label(url, article.get("source_id", "")) # 去重合并后的多来源(如 "Reuters, CNBC"),否则回退单源展示
src = _article_source_label(article)
# 全部来源列表(主源居首,JSON 数组,写入 news_event.sources)
source_list = _article_sources_full(article)
# 归一化:"" / "?" 不入库,留 None(DB 仅存 positive/negative/neutral) # 归一化:"" / "?" 不入库,留 None(DB 仅存 positive/negative/neutral)
sentiment = ev.get("sentiment") or None sentiment = ev.get("sentiment") or None
if sentiment in ("", "?"): if sentiment in ("", "?"):
@@ -594,7 +631,8 @@ def _build_report_data(
title=title or "(无标题)", title=title or "(无标题)",
summary=ev.get("summary_zh") or None, summary=ev.get("summary_zh") or None,
sentiment=sentiment, sentiment=sentiment,
source=src if src not in (None, "", "?") else None, source=src,
sources=source_list,
url=url, url=url,
) )
) )
+74
View File
@@ -0,0 +1,74 @@
#!/bin/bash
# =============================================
# 步骤状态管理(--resume 中断恢复用)
# =============================================
# 用法(在脚本中 source):
# source "$SCRIPT_DIR/_step_state.sh"
# RESUME=0|1 # 由入口脚本参数解析后设置
# step_run <步骤名> <cmd...> # 执行并在成功后标记;resume 时已标记则跳过
# step_should_run <名> # 返回 0 = 需要运行(供需容忍部分失败的步骤手动控制)
# step_mark <名> # 手动标记完成
# =============================================
# 状态文件: data/run_state/{YYYYMMDD}.state(按本地日期换新,跨天自动失效)
# 说明: 各步骤内部的"跳过已处理文件"由各 Python 模块自带(见 pipeline.sh 头注释),
# 本文件仅做"步骤级"断点记录。
# =============================================
STEP_STATE_DIR="data/run_state"
step_state_file() {
echo "${STEP_STATE_DIR}/$(date +%Y%m%d).state"
}
step_done() {
# $1: 步骤名;已标记返回 0
[ -f "$(step_state_file)" ] && grep -qx "$1" "$(step_state_file)" >/dev/null 2>&1
}
step_mark() {
# $1: 步骤名;记录为已完成(幂等追加)
mkdir -p "${STEP_STATE_DIR}"
echo "$1" >> "$(step_state_file)"
}
step_should_run() {
# $1: 步骤名;resume 模式下已标记 → 返回 1(跳过),否则返回 0
if [ "${RESUME:-0}" = "1" ] && step_done "$1"; then
LOG "⏭️ 跳过已完成的步骤: $1(--resume)"
return 1
fi
return 0
}
step_run() {
# $1: 步骤名,$2...: 要执行的命令(成功后标记完成)
# 显式检查退出码:避免 bash set -e 在条件上下文(&&/||)调用函数时失效,
# 导致失败步骤被误标记为完成。
# 设置环境变量 STEP_LOG=路径 时,步骤全部输出同时写入该日志(终端仍实时显示)。
local name="$1"
shift
if ! step_should_run "$name"; then
return 0
fi
local t0 rc t1
t0=$(date +%s)
LOG "▶ $name 开始"
# 显式取原命令退出码(PIPESTATUS[0]):
# - 无 pipefail 时管道码=tee 码(0),必须用 PIPESTATUS 还原真实结果;
# - pipefail + 非条件上下文时失败会触发 set -e 直接退出(此时不 mark,正确);
# - 条件上下文调用时 set -e 被抑制,这里仍能正确捕获并 return。
if [ -n "${STEP_LOG:-}" ]; then
"$@" 2>&1 | tee -a "$STEP_LOG"
rc=${PIPESTATUS[0]}
else
"$@"
rc=$?
fi
t1=$(date +%s)
if [ "$rc" -ne 0 ]; then
LOG "✗ ${name} 失败(exit=${rc},耗时 $((t1 - t0))s);未标记完成,可 --resume 重试"
return "$rc"
fi
LOG "✔ $name 完成(耗时 $((t1 - t0))s)"
step_mark "$name"
}
+38 -5
View File
@@ -6,25 +6,58 @@
# ============================================= # =============================================
# 每天 06:00 / 12:00 / 18:00 / 22:00 各执行一次 # 每天 06:00 / 12:00 / 18:00 / 22:00 各执行一次
# ============================================= # =============================================
# 用法:
# ./scripts/domestic_full.sh # 全新执行
# ./scripts/domestic_full.sh --resume # 从中断处继续(跳过已完成步骤,
# # 步骤状态见 data/run_state/{date}.state)
# =============================================
set -euo pipefail set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
PROJECT_DIR="$(dirname "$SCRIPT_DIR")" PROJECT_DIR="$(dirname "$SCRIPT_DIR")"
cd "$PROJECT_DIR" cd "$PROJECT_DIR"
# ── 参数解析 ──
RESUME=0
for arg in "$@"; do
case "$arg" in
--resume) RESUME=1 ;;
*) echo "未知参数: ${arg}(支持 --resume)" >&2; exit 1 ;;
esac
done
source "$SCRIPT_DIR/_step_state.sh"
LOG() { echo "[$(date '+%Y-%m-%d %H:%M:%S')] $*"; } LOG() { echo "[$(date '+%Y-%m-%d %H:%M:%S')] $*"; }
LOG "══════ 国内全流程开始 ══════" # ── 全流程日志:所有输出实时显示到终端,同时完整写入日志文件 ──
LOG_DIR="logs"
mkdir -p "$LOG_DIR"
FULL_LOG="$LOG_DIR/domestic_full_$(date +%Y%m%d_%H%M%S).log"
: > "$FULL_LOG"
LOG "══════ 国内全流程开始(resume=${RESUME})═══════"
LOG "全流程日志: ${FULL_LOG}(终端实时显示完整进度)"
# ── 加载 .env ── # ── 加载 .env ──
export $(grep -v '^#' .env | grep -v '^$' | xargs 2>/dev/null || true) export $(grep -v '^#' .env | grep -v '^$' | xargs 2>/dev/null || true)
# ── 1. M1 Pi 抓取(headful Playwright + HTTP 代理)── # ── 1. M1 Pi 抓取(headful Playwright + HTTP 代理)──
LOG "[1/2] Pi M1 抓取(8G headful + HTTP 代理)..." # 部分源抓取失败不阻塞管道(原语义);抓取本身按 URL 去重(index.jsonl),
bash "$SCRIPT_DIR/domestic_crawl_8g.sh" 2>&1 | tail -5 || LOG "WARNING: 部分源抓取失败,继续管道" # 已抓取过的 URL 不会重复写入;输出实时显示每源进度
if step_should_run M1_crawl; then
LOG "━━━ [1/2] M1 抓取(headful Playwright + HTTP 代理)━━━"
bash "$SCRIPT_DIR/domestic_crawl_8g.sh" 2>&1 | tee -a "$FULL_LOG" \
|| LOG "WARNING: 部分源抓取失败,继续管道"
step_mark M1_crawl
fi
# ── 2. M2→M6 管道(含日报)── # ── 2. M2→M6 管道(含日报)──
LOG "[2/2] 全链路管道..." LOG "━━━ [2/2] M2→M6 管道 + 日报 ━━━"
bash "$SCRIPT_DIR/pipeline.sh" 2>&1 | tail -10 if [ "$RESUME" = "1" ]; then
bash "$SCRIPT_DIR/pipeline.sh" --resume 2>&1 | tee -a "$FULL_LOG"
else
bash "$SCRIPT_DIR/pipeline.sh" 2>&1 | tee -a "$FULL_LOG"
fi
LOG "══════ 国内全流程完成 ✅ ══════" LOG "══════ 国内全流程完成 ✅ ══════"
+69 -18
View File
@@ -2,9 +2,19 @@
# ============================================= # =============================================
# 国内服务器:全链路管道 M2 → M3 → M4 → M5 → M6 → 日报 # 国内服务器:全链路管道 M2 → M3 → M4 → M5 → M6 → 日报
# ============================================= # =============================================
# 用法:./scripts/pipeline.sh # 用法:
# ./scripts/pipeline.sh # 全新执行(步骤级不跳过)
# ./scripts/pipeline.sh --resume # 从中断处继续(跳过已完成步骤)
# 前提:domestic_sync.sh 已完成 # 前提:domestic_sync.sh 已完成
# ============================================= # =============================================
# 各步骤"跳过已处理文件"(文件级增量,由各 Python 模块自带):
# M2 正文提取: data/processed/{source}/{date}/{url_hash}.json 存在则跳过
# M3 去重: 指纹库 SQLite(data/dedup_fingerprints.sqlite3)判重,重复自动丢弃
# M4 翻译: data/events/{date}/{url_hash}.json 存在则跳过
# M5 向量: data/embeddings/{date}/index.json 记录已向量化,跳过
# M6 入库: Qdrant upsert 幂等(点 id = url_hash,重复写入覆盖,无重复)
# 日报: MySQL 唯一键 (report_date, intl, "") 幂等覆盖
# =============================================
set -euo pipefail set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
@@ -12,57 +22,98 @@ PROJECT_DIR="$(dirname "$SCRIPT_DIR")"
cd "$PROJECT_DIR" cd "$PROJECT_DIR"
echo "[$(date)] ═══ 全链路管道开始 ═══" # ── 参数解析 ──
RESUME=0
for arg in "$@"; do
case "$arg" in
--resume) RESUME=1 ;;
*) echo "未知参数: ${arg}(支持 --resume)" >&2; exit 1 ;;
esac
done
source "$SCRIPT_DIR/_step_state.sh"
LOG() { echo "[$(date '+%Y-%m-%d %H:%M:%S')] $*"; }
# ── 步骤日志:全部输出实时显示到终端,同时完整写入日志文件 ──
LOG_DIR="logs"
mkdir -p "$LOG_DIR"
STEP_LOG="$LOG_DIR/pipeline_$(date +%Y%m%d_%H%M%S).log"
: > "$STEP_LOG"
LOG "══════ 全链路管道开始(resume=${RESUME})═══════"
LOG "步骤日志: ${STEP_LOG}(终端实时显示完整进度)"
# ── AI 模型信息读取(供阶段 banner 显性展示供应商/模型)──
ai_model_info() {
# $1: 场景名(translation / daily_report)或 "embedding"
# 输出格式: "供应商 / 模型"
.venv/bin/python3 -c "
import sys, yaml
scene = sys.argv[1]
raw = yaml.safe_load(open('configs/system.yaml', encoding='utf-8'))
if scene == 'embedding':
cfg = raw.get('embedding', {})
p = cfg.get('provider', '?')
m = cfg.get('dashscope_model', '?')
else:
base = raw.get('llm', {})
sc = (raw.get('llm_scenes', {}) or {}).get(scene, {})
merged = {**base, **sc}
p = merged.get('provider', '?')
m = merged.get('model') or merged.get(p + '_model', '?')
print(f'{p} / {m}')
" "$1"
}
# 加载 .env # 加载 .env
export $(grep -v '^#' .env | grep -v '^$' | xargs 2>/dev/null || true) export $(grep -v '^#' .env | grep -v '^$' | xargs 2>/dev/null || true)
# ── M2: 正文提取 ── # ── M2: 正文提取 ──
echo "[$(date)] [M2] 正文提取..." LOG "━━━ M2 正文提取 ━━━"
.venv/bin/python3 -c " step_run M2_extract .venv/bin/python3 -c "
from extractor.pipeline import process_all_sources from extractor.pipeline import process_all_sources
stats = process_all_sources() stats = process_all_sources()
print(f'M2: {stats[\"total_articles\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s') print(f'M2: {stats[\"total_articles\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s')
" "
# ── M3: 去重 ── # ── M3: 去重 ──
echo "[$(date)] [M3] 三层去重..." LOG "━━━ M3 三层去重 ━━━"
.venv/bin/python3 -c " step_run M3_dedup .venv/bin/python3 -c "
from dedup.pipeline import dedup_all_sources from dedup.pipeline import dedup_all_sources
stats = dedup_all_sources() stats = dedup_all_sources()
print(f'M3: 唯一 {stats[\"unique\"]}/重复 {stats[\"duplicate\"]}, {stats[\"elapsed_sec\"]:.0f}s') print(f'M3: 唯一 {stats[\"unique\"]}/重复 {stats[\"duplicate\"]}, {stats[\"elapsed_sec\"]:.0f}s')
" "
# ── M4: 翻译+事件 ── # ── M4: 翻译+事件(AI 大模型)──
echo "[$(date)] [M4] 翻译+事件抽取..." LOG "━━━ M4 翻译+事件抽取(AI 大模型: $(ai_model_info translation))━━━"
.venv/bin/python3 -c " step_run M4_translate .venv/bin/python3 -c "
from llm.pipeline import translate_all_deduped from llm.pipeline import translate_all_deduped
stats = translate_all_deduped() stats = translate_all_deduped()
print(f'M4: {stats[\"success\"]}/{stats[\"total\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s') print(f'M4: {stats[\"success\"]}/{stats[\"total\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s')
" "
# ── M5: 向量生成 ── # ── M5: 向量生成(AI 大模型)──
echo "[$(date)] [M5] 向量生成..." LOG "━━━ M5 向量生成(AI 大模型: $(ai_model_info embedding))━━━"
.venv/bin/python3 -c " step_run M5_embed .venv/bin/python3 -c "
from embedding.pipeline import embed_all_events from embedding.pipeline import embed_all_events
stats = embed_all_events() stats = embed_all_events()
print(f'M5: {stats[\"success\"]}/{stats[\"total\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s') print(f'M5: {stats[\"success\"]}/{stats[\"total\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s')
" "
# ── M6: Qdrant 入库 ── # ── M6: Qdrant 入库 ──
echo "[$(date)] [M6] Qdrant 入库..." LOG "━━━ M6 Qdrant 入库 ━━━"
.venv/bin/python3 -c " step_run M6_index .venv/bin/python3 -c "
from vectorstore.pipeline import ingest_all_embeddings from vectorstore.pipeline import ingest_all_embeddings
stats = ingest_all_embeddings() stats = ingest_all_embeddings()
print(f'M6: {stats[\"ingested\"]}/{stats[\"total\"]} 条, {stats[\"elapsed_sec\"]:.0f}s') print(f'M6: {stats[\"ingested\"]}/{stats[\"total\"]} 条, {stats[\"elapsed_sec\"]:.0f}s')
" "
# ── 日报 ── # ── 日报(AI 大模型)──
echo "[$(date)] [日报] 生成日报(结构化入库)..." LOG "━━━ 日报生成(AI 大模型: $(ai_model_info daily_report))━━━"
.venv/bin/python3 -c " step_run report .venv/bin/python3 -c "
from scheduler.reporter import generate_report from scheduler.reporter import generate_report
report_id = generate_report() report_id = generate_report()
print(f'日报: report_id={report_id}' if report_id is not None else '日报: 无数据/失败') print(f'日报: report_id={report_id}' if report_id is not None else '日报: 无数据/失败')
" "
echo "[$(date)] ═══ 全链路管道完成 ✅ ═══" LOG "══════ 全链路管道完成 ✅ ═══════"
+99
View File
@@ -1,5 +1,6 @@
"""M3 三层去重模块单元测试。""" """M3 三层去重模块单元测试。"""
import json
from pathlib import Path from pathlib import Path
import pytest import pytest
@@ -17,6 +18,7 @@ from dedup import (
normalize_content, normalize_content,
simhash64, simhash64,
) )
from dedup.pipeline import dedup_source
from extractor.models import ProcessedArticle from extractor.models import ProcessedArticle
# --------------------------------------------------------------------------- # # --------------------------------------------------------------------------- #
@@ -531,3 +533,100 @@ class TestDedupResult:
summary = r.short_summary() summary = r.short_summary()
assert "[DUP/simhash]" in summary assert "[DUP/simhash]" in summary
assert "hd=2" in summary assert "hd=2" in summary
# --------------------------------------------------------------------------- #
# dedup_source 跨源来源合并(M9.2:最终显示新闻记录多个来源)
# --------------------------------------------------------------------------- #
class TestMergeSources:
"""dedup_source 在跨源重复时把来源合并进唯一篇。"""
DATE_STR = "20260805"
def _write_processed(
self,
base: Path,
source_id: str,
article: ProcessedArticle,
) -> None:
"""写入 data/processed/{source_id}/{date}/{url_hash}.json。"""
d = base / "data" / "processed" / source_id / self.DATE_STR
d.mkdir(parents=True, exist_ok=True)
(d / f"{article.url_hash}.json").write_text(
article.model_dump_json(indent=2, ensure_ascii=False), encoding="utf-8"
)
def test_cross_source_merge(self, tmp_path, monkeypatch):
"""同内容两源报道 → 唯一篇 source_ids 记录两个来源,主来源不变。"""
monkeypatch.chdir(tmp_path) # 隔离 data/ 相对路径与默认指纹库
content = ("The Federal Reserve kept interest rates unchanged on Wednesday. "
"Markets rallied in response.")
art_a = _make_article(
url="https://www.reuters.com/business/1", url_hash="aaaa111111111111",
source_id="reuters", content=content,
)
art_b = _make_article(
url="https://www.cnbc.com/2026/1", url_hash="bbbb222222222222",
source_id="cnbc", source_name="CNBC", content=content,
)
self._write_processed(tmp_path, "reuters", art_a)
self._write_processed(tmp_path, "cnbc", art_b)
with Deduper() as deduper:
dedup_source("reuters", deduper, self.DATE_STR)
dedup_source("cnbc", deduper, self.DATE_STR)
# 唯一篇 = reuters(先处理),跨源重复后 source_ids 合并
uniq = tmp_path / "data" / "deduped" / self.DATE_STR / "uniques" / "aaaa111111111111.json"
assert uniq.exists()
data = json.loads(uniq.read_text(encoding="utf-8"))
assert data["source_id"] == "reuters" # 主来源不变
assert data["source_ids"] == ["reuters", "cnbc"]
assert (tmp_path / "data" / "deduped" / self.DATE_STR / "uniques"
/ "bbbb222222222222.json").exists() is False # 重复篇不单独落盘
def test_unique_initializes_source_ids(self, tmp_path, monkeypatch):
"""无重复时唯一篇 source_ids 初始化为 [source_id]。"""
monkeypatch.chdir(tmp_path)
art = _make_article(url="https://x.com/1", url_hash="cccc333333333333",
source_id="ft", source_name="Financial Times")
self._write_processed(tmp_path, "ft", art)
with Deduper() as deduper:
dedup_source("ft", deduper, self.DATE_STR)
uniq = tmp_path / "data" / "deduped" / self.DATE_STR / "uniques" / "cccc333333333333.json"
data = json.loads(uniq.read_text(encoding="utf-8"))
assert data["source_ids"] == ["ft"]
def test_merge_idempotent(self, tmp_path, monkeypatch):
"""同一来源重复出现多次合并时去重(不产生重复来源)。"""
monkeypatch.chdir(tmp_path)
content = "Identical content across sources for idempotent test."
art_a = _make_article(
url="https://www.reuters.com/business/2", url_hash="dddd444444444444",
source_id="reuters", content=content,
)
art_b = _make_article(
url="https://www.cnbc.com/2026/2", url_hash="eeee555555555555",
source_id="cnbc", source_name="CNBC", content=content,
)
art_c = _make_article(
url="https://www.marketwatch.com/2", url_hash="ffff666666666666",
source_id="marketwatch", source_name="MarketWatch", content=content,
)
self._write_processed(tmp_path, "reuters", art_a)
self._write_processed(tmp_path, "cnbc", art_b)
self._write_processed(tmp_path, "marketwatch", art_c)
with Deduper() as deduper:
dedup_source("reuters", deduper, self.DATE_STR)
dedup_source("cnbc", deduper, self.DATE_STR)
dedup_source("marketwatch", deduper, self.DATE_STR)
uniq = tmp_path / "data" / "deduped" / self.DATE_STR / "uniques" / "dddd444444444444.json"
data = json.loads(uniq.read_text(encoding="utf-8"))
assert data["source_ids"] == ["reuters", "cnbc", "marketwatch"]
+6
View File
@@ -87,6 +87,12 @@ class TestLoadEmbeddingConfig:
assert cfg.provider == "dashscope" assert cfg.provider == "dashscope"
assert cfg.api_key == "sk-dashscope-test" assert cfg.api_key == "sk-dashscope-test"
def test_max_attempts_from_system_yaml(self, monkeypatch):
"""重试次数取自 system.yaml embedding.max_attempts(当前 3)。"""
monkeypatch.setenv("DASHSCOPE_API_KEY", "sk-dashscope-test")
cfg = load_embedding_config()
assert cfg.max_attempts == 3
def test_fallback_to_qwen_key(self, monkeypatch): def test_fallback_to_qwen_key(self, monkeypatch):
monkeypatch.delenv("DASHSCOPE_API_KEY", raising=False) monkeypatch.delenv("DASHSCOPE_API_KEY", raising=False)
monkeypatch.setenv("QWEN_API_KEY", "sk-qwen-key") monkeypatch.setenv("QWEN_API_KEY", "sk-qwen-key")
+29
View File
@@ -457,6 +457,35 @@ class TestLoadLLMConfig:
with pytest.raises(ValueError, match="API key 未配置"): with pytest.raises(ValueError, match="API key 未配置"):
load_llm_config(provider="deepseek") load_llm_config(provider="deepseek")
def test_max_attempts_from_system_yaml(self, monkeypatch):
"""重试次数取自 system.yaml llm.max_attempts(当前 3)。"""
monkeypatch.setenv("DEEPSEEK_API_KEY", "sk-deepseek-test-key")
config = load_llm_config(provider="deepseek")
assert config.max_attempts == 3
def test_scene_translation(self, monkeypatch):
"""translation 场景覆盖 llm 默认段(模型/温度)。"""
monkeypatch.setenv("DEEPSEEK_API_KEY", "sk-deepseek-test-key")
config = load_llm_config(provider="deepseek", scene="translation")
assert config.model == "deepseek-v4-flash"
assert config.temperature == 0.1
assert config.max_attempts == 3
def test_scene_daily_report(self, monkeypatch):
"""daily_report 场景独立配置(温度 0.3 / max_tokens 1500)。"""
monkeypatch.setenv("DEEPSEEK_API_KEY", "sk-deepseek-test-key")
config = load_llm_config(scene="daily_report")
assert config.provider == "deepseek"
assert config.model == "deepseek-v4-flash"
assert config.temperature == 0.3
assert config.max_tokens == 1500
def test_scene_unknown_falls_back_to_default(self, monkeypatch):
"""未定义的场景名回退 llm 默认段,不报错。"""
monkeypatch.setenv("DEEPSEEK_API_KEY", "sk-deepseek-test-key")
config = load_llm_config(provider="deepseek", scene="not_exists")
assert config.temperature == 0.1
# --------------------------------------------------------------------------- # # --------------------------------------------------------------------------- #
# translate_and_extract(mock LLM) # translate_and_extract(mock LLM)
+37
View File
@@ -177,3 +177,40 @@ class TestBuildReportData:
Counter(), "") Counter(), "")
# "?" 不写入 DB,留 None # "?" 不写入 DB,留 None
assert r.events[0].sentiment is None assert r.events[0].sentiment is None
def test_multi_source_label(self) -> None:
"""去重合并后的多来源 → source 拼接展示(Reuters, CNBC)+ sources 完整列表。"""
now = datetime(2026, 8, 4, 8, 0, 0)
ev = _fake_high_event("多来源事件", 5, source_id="reuters",
url="https://reuters.com/news/9")
# 模拟 M3 去重合并:source_ids 含两个来源(主源 reuters 居首)
ev["article"]["source_ids"] = ["reuters", "cnbc"]
r = _build_report_data(now, {}, [ev], Counter(), Counter(), Counter(),
Counter(), "")
assert r.events[0].source == "Reuters, CNBC"
assert r.events[0].sources == ["Reuters", "CNBC"] # 主源居首
assert r.events[0].url == "https://reuters.com/news/9"
def test_multi_source_full_list(self) -> None:
"""超过 3 个来源:source 截前 3,sources 存完整列表。"""
now = datetime(2026, 8, 4, 8, 0, 0)
ev = _fake_high_event("四来源事件", 5, source_id="barrons",
url="https://barrons.com/news/1")
ev["article"]["source_ids"] = ["barrons", "cnbc", "reuters", "ft"]
r = _build_report_data(now, {}, [ev], Counter(), Counter(), Counter(),
Counter(), "")
assert r.events[0].source == "Barron's, CNBC, Reuters" # 前 3 拼接
assert r.events[0].sources == [ # 完整 4 个
"Barron's", "CNBC", "Reuters", "Financial Times",
]
def test_single_source_falls_back(self) -> None:
"""source_ids 为空/单一时回退单源逻辑(不拼接)。"""
now = datetime(2026, 8, 4, 8, 0, 0)
ev = _fake_high_event("单来源事件", 4, source_id="investinglive",
url="https://investinglive.com/news/3")
ev["article"]["source_ids"] = ["investinglive"]
r = _build_report_data(now, {}, [ev], Counter(), Counter(), Counter(),
Counter(), "")
assert r.events[0].source == "InvestingLive"
assert r.events[0].sources == ["InvestingLive"]