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
simon d4a55bcaaa feat: M9 日报结构化入库(HTML 改为写入 MySQL news_report/news_event,report_type=intl)
- 新增 report_db/ 包(models/schema/db,复用 news 项目实现,幂等 upsert)
- reporter.py: _build_report_data + generate_report 写库返回 report_id
- pipeline/cli 适配 report_id 返回值;HTML 渲染/上传保留 deprecated
- 新增 tests/test_report_db.py;.env.example 增加 NEWS_DB_* 配置
- 已部署 pi5 并验证:真实生成 report_id=184(12 事件)+ 幂等覆盖
2026-08-04 11:09:17 +08:00
29 changed files with 1373 additions and 117 deletions
+8
View File
@@ -19,3 +19,11 @@ DASHSCOPE_BASE_URL=https://dashscope.aliyuncs.com/api/v1/services/embeddings/tex
# --- Qdrant ---
QDRANT_URL=http://localhost:6333
QDRANT_API_KEY=
# --- 日报结构化入库(与 news 项目共用 MySQL myquant 库) ---
# pi5 通过内网直连 pi(192.168.1.10) 上的 autossh 隧道(0.0.0.0:13306 → doorcome.cn:3306)
NEWS_DB_HOST=192.168.1.10
NEWS_DB_PORT=13306
NEWS_DB_USER=myquant
NEWS_DB_PASSWORD= # 填真实值(与 news 项目 .env 的 NEWS_DB_PASSWORD 相同),禁止写入源码
NEWS_DB_NAME=myquant
+2
View File
@@ -27,10 +27,12 @@ Thumbs.db
data/raw/
data/processed/
data/deduped/
data/dedup/
data/events/
data/embeddings/
data/qdrant_storage/
data/reports/
data/run_state/
# Logs
logs/*.log
+45 -21
View File
@@ -2,9 +2,38 @@
# Claude Code 开发约束(English Financial News 项目)
版本:v1.0
版本:v1.1
最后更新:2026-06-21
最后更新:2026-07-23
---
# 〇、命令速查(2026-07-23 实测)
开发环境:
* Python 3.11 + uv + `.venv`(详见第三节)
* 安装依赖:`uv sync`;新增依赖:`uv add 包名`;开发依赖:`uv add --dev 包名`
常用命令:
| 命令 | 用途 |
|------|------|
| `uv run pytest` | 跑全部测试(当前 162 passed / 2 failed,见下方已知问题) |
| `uv run ruff check .` | 代码规范检查(当前 13 错误,9 可自动修复) |
| `uv run en-news pipeline` | M2→M6 完整管道 |
| `uv run en-news report` | 仅生成日报 |
| `uv run en-news search "查询"` | Qdrant 语义检索 |
| `uv run en-news mcp-server` | 启动 M8 MCP 服务 |
| `bash scripts/domestic_full.sh` | 全流程(M1→M6→日报),Pi 上由 crontab 06/12/18/22 调度 |
| `bash scripts/domestic_crawl_8g.sh [源ID]` | 仅 M1 抓取(可单源测试) |
| `bash scripts/pipeline.sh` | M2→M6 管道 |
已知问题:
* `tests/test_crawler.py::test_write_and_load_index_jsonl`、`test_write_index_jsonl_dedup` 失败 — `crawler/storage.py:load_index()` 硬编码 `data/raw/...` 路径,未使用测试临时目录;
* ruff 未清零(多为 import 排序 I001,可 `--fix`);
* 日报/摘要 Prompt 内嵌在 `scheduler/reporter.py`,尚未拆分到 `prompts/`。
---
@@ -16,7 +45,7 @@
核心能力包括:
* 英文财经新闻抓取(Crawl4AI,12 个英文源,部署海外);
* 英文财经新闻抓取(Crawl4AI,12 个活跃源,国内 Pi 独立运行);
* 英文正文提取(trafilatura);
* 全文英译中(LLM DeepSeek);
* 投资事件抽取(LLM);
@@ -99,25 +128,20 @@ M0 → M1 → 同步机制 → M2 → M3 → M4 → M5+M6 → M7+日报 → M8
## 2.4 部署拓扑意识
本项目采用双服务器架构(详见 english-news-plan.md 第二节):
本项目当前为单服务器架构(海外服务器已于 2026-07-14 停用,详见 README):
| 服务器 | SSH | 部署路径 | 职责 |
|--------|-----|---------|------|
| 海外 | `ssh ecs-user@8.217.19.253` | `/opt/intlgrab` | M1 抓取 → rsync 推送 |
| 国内 | `ssh pi@192.168.1.160` | `/home/pi/intlnews` | M2→M8 全链路 + 日报 |
| 国内 | `ssh pi@192.168.1.160` | `/home/pi/intlnews` | M1→M8 全链路 + 日报 |
部署优先级:
* 抓取代码优先在海外服务器测试;
* 其余代码在国内服务器测试;
* 海外验证通过后,抓取代码同步到国内做全链路联调。
历史:海外 `ecs-user@8.217.19.253:/opt/intlgrab` 曾负责 M1 抓取 → rsync 推送,已停用;`scripts/overseas_*.sh`、`scripts/domestic_sync.sh` 为遗留脚本(未清理)。
编码时注意:
* 海外侧只负责抓取,无 LLM/Embedding 依赖;
* 国内侧等待 rsync 同步完成后再启动管道;
* 哨兵文件(`SYNC_SENTINEL`)是同步完整性的关键信号;
* `data/raw/` 目录在两台服务器上路径一致。
* M1 抓取在国内 Pi(8GB)以 headful Playwright + HTTP 代理(Shadowsocks + Privoxy)运行,需 Xvfb 虚拟显示器(见 README 前置条件);
* 强反爬源(Reuters / Investing.com / FT)走 Google News RSS 或 RSS 摘要,全文抓取成功率取决于代理线路质量;
* 抓取配置按运行环境拆分为 `configs/profiles/2g_headless.yaml` 与 `configs/profiles/8g_headful.yaml`,通过环境变量 `EN_NEWS_PROFILE` 选择;
* `data/raw/{source_id}/{yyyymmdd}/` 存放抓取原始 HTML + Markdown + index.jsonl,为全管道数据源头。
---
@@ -348,8 +372,7 @@ test_xxx.py
* crawler(M1 抓取)
* extractor(M2 正文提取)
* dedup(M3 去重)
* translator(M4a 翻译)
* llm(M4b 事件抽取)
* llm(M4a 翻译 + M4b 事件抽取,单次 LLM 调用合并输出;`translator/` 为空壳目录,翻译实现在 `llm/`)
* embedding(M5 向量生成)
* vectorstore(M6 Qdrant 入库/检索)
* scheduler(M7 定时调度)
@@ -428,8 +451,8 @@ prompts/
本项目 Prompt 文件:
* `translation_and_extraction.md` — 翻译 + 事件抽取(单次 LLM 调用合并输出)
* `daily_report.md` — 日报生成 Prompt(五段式结构)
* `search_agent.md` — MCP Agent 深度研究 Prompt
注意:`daily_report.md`、`search_agent.md` 尚未创建;日报/摘要 Prompt 目前内嵌在 `scheduler/reporter.py`(技术债,待拆分到 prompts/)。
---
@@ -469,13 +492,14 @@ docker compose up -d
注意:
* 海外侧 docker-compose 仅包含 Crawl4AI + rsync daemon;
* 国内侧 docker-compose 包含 Qdrant + 调度器 + MCP 服务。
* 预期形态:docker-compose 包含 Qdrant + 调度器 + MCP 服务(尚未落地)。
不得依赖:
手工安装。
当前状态(2026-07-23 核查):仓库中尚无 Dockerfile / docker-compose.yml,此要求尚未落实,属待办事项。
---
# 十二、Git 提交规范
+65 -7
View File
@@ -10,7 +10,7 @@
- **M4** 全文英译中 + 投资事件抽取(LLM DeepSeek v4-flash)
- **M5** 向量生成(DashScope text-embedding-v3,1024 维)
- **M6** 双语向量知识库(Qdrant)
- **M7** 定时调度(06/12/18/22) + 每日 AI 摘要日报
- **M7** 定时调度(06/12/18/22) + 每日 AI 摘要日报(M9 起结构化写入 MySQL,不再产出 HTML)
- **M8** MCP 服务(Cherry Studio / Claude Code Agent 深度研究)
## 部署架构
@@ -27,7 +27,7 @@
│ ↓ │
│ M5 向量生成 → M6 Qdrant 入库 │
│ ↓ │
│ 日报生成 → rsync 上传 doorcome.cn │
│ 日报生成 → 写入 MySQL(news_report) │
│ ↓ │
│ M8 MCP 服务(研究 Agent) │
└──────────────────────────────────────┘
@@ -104,6 +104,12 @@ sudo systemctl enable --now privoxy
# 全流程:M1 抓取 → M2→M6 管道 → 日报
# 每天 06:00 / 12:00 / 18:00 / 22:00 自动执行
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
```
### 单步执行
@@ -119,7 +125,7 @@ bash scripts/domestic_crawl_8g.sh investinglive
# M2→M6 管道(需先有 raw 数据)
bash scripts/pipeline.sh
# 仅生成日报
# 仅生成日报(写入 MySQL news_report/news_event,report_type="intl")
uv run en-news report
```
@@ -140,11 +146,63 @@ ls -lt data/reports/
| 文件 | 用途 |
|------|------|
| `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 抓取配置(代理/超时) |
| `.env` | 密钥 / 服务地址(不入 Git) |
| `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)
日报内容结构化写入与 [news 项目](https://github.com/) 共用的 MySQL `myquant` 库(表 `news_report` / `news_event`,`report_type="intl"`,同一天重复生成幂等覆盖)。表结构与数据契约见 news 项目 `docs/db_schema.md`。
`.env` 需配置(与 news 项目 `.env` 的 `NEWS_DB_PASSWORD` 相同):
```env
NEWS_DB_HOST=192.168.1.10 # pi5 经内网直连 pi 上的 autossh 隧道(0.0.0.0:13306 → doorcome.cn:3306)
NEWS_DB_PORT=13306
NEWS_DB_USER=myquant
NEWS_DB_PASSWORD=xxx
NEWS_DB_NAME=myquant
```
---
## 数据存储
@@ -156,7 +214,7 @@ data/
├── dedup/ M3 去重指纹库(SQLite)
├── translations/{yyyymmdd}/ M4 翻译+事件抽取结果(JSONL)
├── embeddings/{yyyymmdd}/ M5 向量文件(JSONL)
├── reports/ M7 日报(HTML)
├── reports/ M7 日报(M9 起不再产出 HTML,改存 MySQL)
└── vectorstore/ M6 Qdrant 数据
```
@@ -191,9 +249,9 @@ uv run en-news mcp-server
部分网站有强反爬(Cloudflare/DataDome/付费墙),RSS 只返回摘要。全文可尝试 headful 浏览器回退,但成功率取决于代理线路质量。
### 日报上传到哪里?
### 日报存在哪里?
日报自动上传到 `simon@doorcome.cn:/var/www/html/echart/research/`,通过 HTTP 访问。
M9 起日报结构化写入 MySQL `myquant` 库:主表 `news_report`(`report_type="intl"`)+ 明细表 `news_event`(`section="intl"`),数据契约与 [news 项目](https://github.com/) 一致(见其 `docs/db_schema.md`)。历史 HTML 日报(2026-06-16 ~ 2026-08-03)已由 news 项目解析导入。API / 前端读取同一批表。
### 如何添加新源?
+7 -6
View File
@@ -3,8 +3,8 @@
import logging
import sys
from dotenv import load_dotenv
import typer
from dotenv import load_dotenv
# 加载 .env 环境变量(必须在所有模块导入之前)
load_dotenv()
@@ -44,6 +44,7 @@ def crawl(
):
"""M1: 抓取英文财经新闻(部署海外)"""
import os
from crawler.orchestrator import run_crawl_sync
if profile:
@@ -196,15 +197,15 @@ def search(
@app.command()
def report():
"""M7: 日报生成"""
"""M7: 日报生成(结构化入库,不再产出 HTML)"""
from scheduler.reporter import generate_report
try:
path = generate_report()
if path:
typer.echo(f"\n✅ 日报已生成: {path}")
report_id = generate_report()
if report_id is not None:
typer.echo(f"\n✅ 日报已入库: report_id={report_id}")
else:
typer.echo("⚠️ 无数据,跳过日报生成")
typer.echo("⚠️ 无数据或入库失败,跳过日报生成")
except Exception as e:
logger.exception("日报生成失败")
typer.echo(f"❌ 日报生成出错: {e}", err=True)
+43 -4
View File
@@ -40,23 +40,62 @@ dedup:
simhash_window_days: 30
min_content_length: 100
# ── LLM 翻译+事件抽取 ────────────────────────────────
# ── LLM 默认配置(所有 LLM 场景的兜底)────────────────
# 按场景独立配置见下方 llm_scenes 段;场景未声明的字段回退到本段。
llm:
provider: "deepseek"
deepseek_model: "deepseek-chat"
deepseek_model: "deepseek-v4-flash"
qwen_model: "qwen-plus"
timeout_sec: 60
max_retries: 3
max_attempts: 3 # 单篇总尝试次数(含首次),失败后指数退避重试
max_tokens: 8192
temperature: 0.1
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:
provider: "dashscope"
dashscope_model: "text-embedding-v3"
dimension: 1024
batch_size: 10
max_attempts: 3 # 单批总尝试次数(含首次),失败后指数退避重试
timeout_sec: 30
# ── Qdrant ──────────────────────────────────────────
+137 -1
View File
@@ -1,6 +1,142 @@
# continuation.md — English Financial News 项目状态
> 最后更新:2026-07-23
> 最后更新: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
---
## 2026-08-04 会话成果
### M9:日报结构化入库(HTML → MySQL)
**背景:** 与 news 项目对齐,日报内容不再产出 HTML / scp 上传 doorcome,改为写入共用 MySQL `myquant` 库(表 `news_report` / `news_event`,`report_type="intl"`)。历史日报(2026-06-16 ~ 2026-08-03,intl 128 份)已由 news 项目解析入库,本次仅做新日报按相同契约入库。
**改动:**
| 文件 | 改动内容 |
|------|---------|
| `report_db/`(新增) | 复用 news 项目实现:`models.py`(EventRow/ReportData)、`schema.py`(news_report/news_event DDL)、`db.py`(load_db_config/connect/init_schema/save_report/fetch_report/exists_report,loguru 换 logging) |
| `scheduler/reporter.py` | 新增 `_build_report_data()`(intl 结构:事件全部 `section="intl"`,`file_name=""` 幂等覆盖,stats 含 pipeline/sentiment/importance/event_types/source_dist);`generate_report()` 返回 `int | None`,收集逻辑不变,末尾 `connect()+save_report()`;HTML 渲染/上传函数保留标 deprecated |
| `scheduler/pipeline.py` | `run_step_report` 适配 `report_id` 返回值 |
| `app/cli.py` | `report` 命令打印 `report_id` |
| `pyproject.toml` | `uv add pymysql`;`report_db/schema.py` 加 E501 per-file-ignore |
| `.env.example` / `README.md` | 新增 `NEWS_DB_*` 配置说明(pi5 经内网直连 `192.168.1.10:13306`,password 与 news 项目相同) |
| `tests/test_report_db.py`(新增) | 模型默认值 + `_build_report_data` 组装(板块/rank、字段映射、标题截断、stats 快照、`"?"` 归一化) |
**测试结果:** `uv run pytest` → 174 passed / 2 failed(2 个失败为原有 `test_crawler.py` index.jsonl 路径问题,与本次无关);ruff 改动文件干净(仅剩原有 `_SENTENCE_END` N806)。
**部署:** 已 rsync 至 pi5(`/home/pi/intlnews`),`.env` 写入 NEWS_DB_*,`uv sync` 安装 pymysql,真实生成日报并回读 DB 验证(详见下方)。
---
+37
View File
@@ -103,6 +103,8 @@ def dedup_source(
if result.is_duplicate:
dup_count += 1
# 跨源重复:把来源合并进已保留的唯一篇(记录多个来源)
_merge_duplicate_source(article, result)
logger.debug("[%s] 🔁 %s → L%d: %s",
source_id,
article.title[:40],
@@ -110,6 +112,9 @@ def dedup_source(
result.short_summary())
else:
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.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:
"""DedupResult → 命中层编号。"""
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))
batch_size = min(int(sys_cfg.get("batch_size", DASHSCOPE_BATCH_LIMIT)), DASHSCOPE_BATCH_LIMIT)
timeout = float(sys_cfg.get("timeout_sec", 30.0))
max_attempts = int(sys_cfg.get("max_attempts", DEFAULT_MAX_ATTEMPTS))
if not api_key:
raise EmbeddingError("DASHSCOPE_API_KEY 未配置,请检查 .env")
@@ -104,6 +105,7 @@ def load_embedding_config(
dimension=dimension,
batch_size=batch_size,
timeout_sec=timeout,
max_attempts=max_attempts,
)
+2 -2
View File
@@ -100,8 +100,8 @@ def embed_all_events(
# 批量嵌入(按 batch_size 分块,每批输出进度)
batch_size = config.batch_size
total = len(articles)
logger.info("开始向量化 %d 篇文章(batch_size=%d, model=%s)",
total, batch_size, config.model)
logger.info("开始向量化 %d 篇文章(batch_size=%d, provider=%s, model=%s)",
total, batch_size, config.provider, config.model)
for batch_start in range(0, total, batch_size):
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):
@@ -20,3 +20,8 @@ class ProcessedArticle(BaseModel):
status: str = "success" # success | no_content | failed
extractor: str = "trafilatura" # trafilatura | crawl4ai_md | none
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"
_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")
if config_path.exists():
try:
with open(config_path, encoding="utf-8") as 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:
logger.warning("加载 llm 配置失败,使用空配置")
return {}
@@ -50,6 +65,7 @@ class LLMConfig:
timeout_sec: float = 60.0
temperature: float = 0.1
max_tokens: int = 8192
max_attempts: int = _DEFAULT_MAX_ATTEMPTS # 单次调用总尝试次数(含首次)
def __post_init__(self) -> None:
if not self.api_key:
@@ -60,26 +76,28 @@ def load_llm_config(
provider: str | None = None,
*,
model: str | None = None,
scene: str | None = None,
) -> LLMConfig:
"""根据配置文件构造 LLMConfig。
provider 为 None 时读 system.yaml llm.provider,默认 deepseek。
model 为 None 时读 system.yaml 中对应 provider 的 model。
provider 为 None 时读 system.yaml(或场景段)llm.provider,默认 deepseek。
model 为 None 时优先读场景段的 model 键,其次 system.yaml 对应 provider 的 model。
scene 指定时,llm_scenes.{scene} 覆盖 llm 默认段的各参数(按场景独立配置模型)。
Raises:
ValueError: API key 未配置
"""
config = _load_system_config()
config = _load_system_config(scene=scene)
p = (provider or config.get("provider", "deepseek")).lower()
if p == "deepseek":
api_key = os.environ.get("DEEPSEEK_API_KEY", "")
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"):
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)
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"
else:
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))
temperature = float(config.get("temperature", 0.1))
max_tokens = int(config.get("max_tokens", 8192))
max_attempts = int(config.get("max_attempts", _DEFAULT_MAX_ATTEMPTS))
return LLMConfig(
provider=p,
@@ -102,6 +121,7 @@ def load_llm_config(
timeout_sec=timeout,
temperature=temperature,
max_tokens=max_tokens,
max_attempts=max_attempts,
)
+10 -3
View File
@@ -226,7 +226,7 @@ def translate_and_extract(
article: ProcessedArticle,
*,
template: PromptTemplate | None = None,
max_attempts: int = DEFAULT_MAX_ATTEMPTS,
max_attempts: int | None = None,
) -> EnTranslatedArticle:
"""同步翻译 + 事件抽取(单篇文章,带重试)。
@@ -235,7 +235,8 @@ def translate_and_extract(
config: LLM 配置
article: 待处理的英文新闻
template: Prompt 模板,默认加载 prompts/translation_and_extraction.md
max_attempts: 最大重试次数
max_attempts: 最大尝试次数(含首次);None 时取 config.max_attempts
(来自 system.yaml llm.max_attempts)
Returns:
EnTranslatedArticle 含双语内容 + 事件
@@ -243,6 +244,8 @@ def translate_and_extract(
Raises:
LLMCallError: 所有重试均失败
"""
if max_attempts is None:
max_attempts = config.max_attempts
tpl = template or PromptTemplate()
system_prompt, user_prompt = tpl.render(article)
@@ -261,6 +264,7 @@ def translate_and_extract(
source_name=article.source_name,
url=article.url,
url_hash=article.url_hash,
source_ids=list(article.source_ids),
title=article.title,
title_zh=output.title_zh,
content_en=article.content,
@@ -305,10 +309,12 @@ async def translate_and_extract_async(
article: ProcessedArticle,
*,
template: PromptTemplate | None = None,
max_attempts: int = DEFAULT_MAX_ATTEMPTS,
max_attempts: int | None = None,
semaphore: asyncio.Semaphore | None = None,
) -> EnTranslatedArticle:
"""异步翻译 + 事件抽取(批处理用),与同步版逻辑等价。"""
if max_attempts is None:
max_attempts = config.max_attempts
tpl = template or PromptTemplate()
system_prompt, user_prompt = tpl.render(article)
@@ -327,6 +333,7 @@ async def translate_and_extract_async(
source_name=article.source_name,
url=article.url,
url_hash=article.url_hash,
source_ids=list(article.source_ids),
title=article.title,
title_zh=output.title_zh,
content_en=article.content,
+4
View File
@@ -103,6 +103,10 @@ class EnTranslatedArticle(BaseModel):
source_name: str
url: str
url_hash: str
source_ids: list[str] = Field(
default_factory=list,
description="去重合并后的所有来源 ID(透传自 ProcessedArticle.source_ids)",
)
# ── 双语内容 ──
title: str = "" # 英文原标题
+6 -2
View File
@@ -119,9 +119,13 @@ def translate_all_deduped(
logger.warning("去重目录无文章: data/deduped/%s/uniques/", date_str)
return {"date": date_str, "total": 0, "success": 0, "failed": 0, "elapsed_sec": 0}
# 初始化 LLM 客户端
config = load_llm_config(provider=provider, model=model)
# 初始化 LLM 客户端(translation 场景配置见 system.yaml llm_scenes.translation)
config = load_llm_config(provider=provider, model=model, scene="translation")
client = make_sync_client(config)
logger.info(
"AI 大模型(场景 translation): provider=%s model=%s",
config.provider, config.model,
)
template = PromptTemplate()
# 输出目录
+2
View File
@@ -25,6 +25,7 @@ dependencies = [
"openai>=2.43.0",
"python-dotenv>=1.2.2",
"markdown>=3.10.2",
"pymysql>=1.2.0",
]
[project.optional-dependencies]
@@ -65,6 +66,7 @@ select = ["E", "F", "I", "N", "W", "UP"]
[tool.ruff.lint.per-file-ignores]
"scheduler/reporter.py" = ["E501"] # 内嵌 HTML 模板 CSS 行较长
"report_db/schema.py" = ["E501"] # DDL SQL 语句较长,与 news 项目 schema 保持一致
[tool.pytest.ini_options]
testpaths = ["tests"]
+20
View File
@@ -0,0 +1,20 @@
"""日报结构化入库。
职责:日报内容(intl)结构化后写入 MySQL(news_report / news_event),
与 news 项目共用同一表结构与数据契约(见 docs/db_schema.md、report_db_design.md)。
本包不提供 API 与前端。
"""
from .db import connect, exists_report, fetch_report, init_schema, load_db_config, save_report
from .models import EventRow, ReportData
__all__ = [
"EventRow",
"ReportData",
"connect",
"exists_report",
"fetch_report",
"init_schema",
"load_db_config",
"save_report",
]
+163
View File
@@ -0,0 +1,163 @@
"""日报结构化入库:DB 连接、建表、写入。
与 news 项目 report_db/db.py 保持一致(日志改用 logging,遵守本项目规范)。
"""
from __future__ import annotations
import json
import logging
import os
from dataclasses import dataclass
from typing import Any
import pymysql
from .models import ReportData
from .schema import DDL_STATEMENTS
logger = logging.getLogger(__name__)
@dataclass(frozen=True)
class DbConfig:
"""MySQL 连接配置(来自环境变量 NEWS_DB_*)。"""
host: str
port: int
user: str
password: str
name: str
def load_db_config() -> DbConfig:
"""从环境变量读取 NEWS_DB_*,缺失密码时抛异常(禁止默认密码)。"""
host = os.environ.get("NEWS_DB_HOST", "192.168.1.10")
port = int(os.environ.get("NEWS_DB_PORT", "13306"))
user = os.environ.get("NEWS_DB_USER", "myquant")
password = os.environ.get("NEWS_DB_PASSWORD", "")
name = os.environ.get("NEWS_DB_NAME", "myquant")
if not password:
logger.error("NEWS_DB_PASSWORD 未配置,请在 .env 中设置")
raise ValueError("NEWS_DB_PASSWORD 未配置")
return DbConfig(host=host, port=port, user=user, password=password, name=name)
def connect(cfg: DbConfig | None = None) -> pymysql.Connection:
"""建立短连接(autocommit=False)。失败时记录日志并抛出。"""
cfg = cfg or load_db_config()
try:
conn = pymysql.connect(
host=cfg.host,
port=cfg.port,
user=cfg.user,
password=cfg.password,
database=cfg.name,
charset="utf8mb4",
autocommit=False,
cursorclass=pymysql.cursors.DictCursor,
)
except Exception:
logger.exception("连接 MySQL 失败: host=%s port=%s user=%s", cfg.host, cfg.port, cfg.user)
raise
logger.debug("MySQL 已连接: %s/%s", cfg.host, cfg.name)
return conn
def init_schema(conn: pymysql.Connection) -> None:
"""建表(CREATE TABLE IF NOT EXISTS ×2),幂等。"""
with conn.cursor() as cur:
for ddl in DDL_STATEMENTS:
cur.execute(ddl)
conn.commit()
logger.info("news_report / news_event 建表完成")
def save_report(conn: pymysql.Connection, report: ReportData) -> int:
"""事务内写入一份日报。
幂等策略:
- 主表按 (report_date, report_type, file_name) 唯一键 upsert;
- 事件表 DELETE 该 report 旧行后全量 INSERT(整份覆盖一致)。
返回 report_id。
"""
with conn.cursor() as cur:
stats_json = json.dumps(report.stats, ensure_ascii=False) if report.stats else None
cur.execute(
"""
INSERT INTO news_report
(report_date, report_type, file_name, generated_at, ai_summary, stats)
VALUES (%s, %s, %s, %s, %s, %s)
ON DUPLICATE KEY UPDATE
generated_at = VALUES(generated_at),
ai_summary = VALUES(ai_summary),
stats = VALUES(stats)
""",
(
report.report_date,
report.report_type,
report.file_name,
report.generated_at,
report.ai_summary,
stats_json,
),
)
cur.execute(
"SELECT id FROM news_report WHERE report_date=%s AND report_type=%s AND file_name=%s",
(report.report_date, report.report_type, report.file_name),
)
row = cur.fetchone()
if row is None: # pragma: no cover - 理论不可达
raise RuntimeError("写入 news_report 后查询不到 report_id")
report_id: int = row["id"]
cur.execute("DELETE FROM news_event WHERE report_id=%s", (report_id,))
for ev in report.events:
sources_json = json.dumps(ev.sources, ensure_ascii=False) if ev.sources else None
cur.execute(
"""
INSERT INTO news_event
(report_id, section, rank, importance, event_type, title,
summary, sentiment, source, sources, url)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
""",
(
report_id,
ev.section,
ev.rank,
ev.importance,
ev.event_type,
ev.title,
ev.summary,
ev.sentiment,
ev.source,
sources_json,
ev.url,
),
)
conn.commit()
logger.info("日报已入库: report_id=%s date=%s type=%s events=%s",
report_id, report.report_date, report.report_type, len(report.events))
return report_id
def fetch_report(conn: pymysql.Connection, report_id: int) -> dict[str, Any] | None:
"""读侧辅助(联调/测试用),返回主表行。"""
with conn.cursor() as cur:
cur.execute("SELECT * FROM news_report WHERE id=%s", (report_id,))
return cur.fetchone()
def exists_report(
conn: pymysql.Connection,
report_date: Any,
report_type: str,
file_name: str,
) -> bool:
"""判断主表是否已存在该唯一键记录(历史导入幂等用)。"""
with conn.cursor() as cur:
cur.execute(
"SELECT 1 FROM news_report WHERE report_date=%s AND report_type=%s AND file_name=%s",
(report_date, report_type, file_name),
)
return cur.fetchone() is not None
+41
View File
@@ -0,0 +1,41 @@
"""日报结构化入库:数据模型。
与 news 项目 report_db/models.py 保持一致(intl 板块复用同一表结构)。
"""
from __future__ import annotations
from datetime import date, datetime
from typing import Any
from pydantic import BaseModel, Field
class EventRow(BaseModel):
"""一条事件记录,对应 news_event 一行。"""
section: str # xwlb | news | cninfo | intl
rank: int # 板块内序号(从 1 开始)
importance: int | None = None
event_type: str | None = None
title: str
summary: str | None = None
sentiment: str | None = None # positive | negative | neutral | ''
source: str | None = None
sources: list[str] = Field(
default_factory=list,
description="全部来源展示名(JSON 数组,主源居首;多源新闻记录完整列表)",
)
url: str | None = None
class ReportData(BaseModel):
"""一份完整日报:news_report 一行 + news_event 多行。"""
report_date: date
report_type: str # finance | intl
file_name: str = "" # 源文件名(新生成日报可为空)
generated_at: datetime
ai_summary: str | None = None
stats: dict[str, Any] = Field(default_factory=dict)
events: list[EventRow] = Field(default_factory=list)
+51
View File
@@ -0,0 +1,51 @@
"""MySQL DDL:日报结构化入库(表前缀 news_,目标 MariaDB 10.11)。
与 news 项目 docs/db_schema.md 保持一致(news_report / news_event)。
历史日报已由 news 项目 report-import 导入,本包只负责新日报写入与幂等建表。
"""
from __future__ import annotations
DDL_STATEMENTS: list[str] = [
# 日报主表
"""
CREATE TABLE IF NOT EXISTS news_report (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
report_date DATE NOT NULL COMMENT '日报日期',
report_type VARCHAR(16) NOT NULL COMMENT 'finance=A股日报 / intl=国际财经日报',
file_name VARCHAR(160) NOT NULL DEFAULT '' COMMENT '源文件名(历史解析);新生成可为空',
generated_at DATETIME NOT NULL COMMENT '生成时间',
ai_summary TEXT NULL COMMENT 'AI 摘要全文',
stats JSON NULL COMMENT '数据总览统计快照(管道/情绪/重要度/事件类型/来源分布)',
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_report_file (report_date, report_type, file_name),
KEY idx_report_date (report_date)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='每日日报主表'
""",
# 日报事件明细
"""
CREATE TABLE IF NOT EXISTS news_event (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
report_id BIGINT UNSIGNED NOT NULL COMMENT 'FK → news_report.id',
section VARCHAR(16) NOT NULL COMMENT '板块: xwlb=新闻联播 / news=财经新闻 / cninfo=公告调研 / intl=国际重要事件',
rank INT NOT NULL DEFAULT 0 COMMENT '板块内序号',
importance INT NULL COMMENT '重要度 1-5',
event_type VARCHAR(64) NULL COMMENT '事件类型',
title VARCHAR(512) NOT NULL COMMENT '标题',
summary TEXT NULL COMMENT '摘要/正文',
sentiment VARCHAR(8) NULL COMMENT 'positive/negative/neutral',
source VARCHAR(64) NULL COMMENT '来源(多源时拼接展示名,最多前 3 个)',
sources TEXT NULL COMMENT '全部来源(JSON 数组,主源居首;多源新闻完整列表)',
url VARCHAR(512) NULL COMMENT '原文链接',
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
KEY idx_report_section (report_id, section),
KEY idx_title (title(255))
) 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
""",
]
+5 -4
View File
@@ -150,16 +150,17 @@ def run_step_index(date_str: str) -> StepResult:
def run_step_report(date_str: str) -> StepResult:
"""日报生成。"""
"""日报生成(M9:结构化入库,不再产出 HTML)。"""
started = datetime.now()
try:
from scheduler.reporter import generate_report
path = generate_report()
report_id = generate_report()
elapsed = (datetime.now() - started).total_seconds()
ok = path is not None
ok = report_id is not None
return StepResult(
name="report", success=ok, elapsed_sec=elapsed,
message=str(path) if path else "无数据", started_at=started,
message=f"report_id={report_id}" if report_id is not None else "无数据",
started_at=started,
)
except Exception as e:
elapsed = (datetime.now() - started).total_seconds()
+156 -33
View File
@@ -9,12 +9,15 @@
import json
import logging
from collections import Counter
from datetime import datetime, timedelta, timezone
from datetime import UTC, datetime, timedelta
from pathlib import Path
from typing import Any
import markdown
import yaml
from report_db.models import EventRow, ReportData # noqa: F401 - 供 _build_report_data 注解使用
logger = logging.getLogger(__name__)
_MAX_HIGH_EVENTS = 30
@@ -93,6 +96,33 @@ def _url_source_label(url: str, source_id: str = "") -> str:
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
@@ -130,8 +160,8 @@ def _try_parse_time(time_str: str) -> datetime | None:
# - 有时区 → astimezone 转为 UTC
# - 无时区 → 假设为 UTC(多数财经新闻 API 使用 UTC)
if dt.tzinfo is None:
return dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc)
return dt.replace(tzinfo=UTC)
return dt.astimezone(UTC)
except ValueError:
continue
return None
@@ -168,7 +198,7 @@ def _load_events_window(now: datetime) -> list[dict]:
文章 dict 列表,含 title/title_zh/url/source_id/events 等
"""
# cutoff 使用 UTC-aware,与 _try_parse_time 返回的 UTC datetime 对齐
cutoff = now.astimezone(timezone.utc) - timedelta(hours=_REPORT_WINDOW_HOURS)
cutoff = now.astimezone(UTC) - timedelta(hours=_REPORT_WINDOW_HOURS)
articles: list[dict] = []
skipped_empty_pt = 0
for day_str in _dates_in_window(now):
@@ -280,7 +310,7 @@ def _call_llm_simple(
system_prompt: str,
user_prompt: str,
max_tokens: int = 600,
max_retries: int = 2,
max_retries: int | None = None,
) -> str:
"""封装 LLM 调用,带重试和客户端复用。
@@ -288,23 +318,30 @@ def _call_llm_simple(
system_prompt: system role 内容
user_prompt: user role 内容
max_tokens: 最大输出 token
max_retries: 最大重试次数(不含首次调用)
max_retries: 重试次数(不含首次);None 时取
config.max_attempts - 1(system.yaml llm.max_attempts)
Returns:
LLM 输出文本;所有重试均失败返回空字符串
"""
import time as _time
# 复用客户端(同 provider/model 只创建一次)
# 复用客户端(同 provider/model 只创建一次;日报摘要场景见 system.yaml llm_scenes.daily_report)
cache_key = "default"
if cache_key not in _llm_client_cache:
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"])
logger.info("AI 大模型(场景 daily_report): provider=%s model=%s",
_llm_client_cache["config"].provider,
_llm_client_cache["config"].model)
config = _llm_client_cache["config"]
client = _llm_client_cache[cache_key]
if max_retries is None:
max_retries = max(config.max_attempts - 1, 0)
last_err: str = ""
for attempt in range(1, max_retries + 2): # 首次 + max_retries 次重试
try:
@@ -314,7 +351,7 @@ def _call_llm_simple(
{"role": "system", "content": system_prompt},
{"role": "user", "content": user_prompt},
],
temperature=0.3,
temperature=config.temperature,
max_tokens=max_tokens,
)
content = (resp.choices[0].message.content or "").strip()
@@ -551,17 +588,102 @@ def _generate_ai_summary(articles: list[dict], day_str: str) -> str:
return "⚠️ AI 摘要合并失败,以下为各批次原始摘要:\n\n" + "\n\n".join(partial_summaries)
def generate_report() -> Path | None:
"""生成 HTML 日报(覆盖过去 25 小时数据)。
def _build_report_data(
now: datetime,
stats: dict,
high_events: list[dict],
sentiments: Counter,
importances: Counter,
event_types: Counter,
sources: Counter,
ai_summary: str,
) -> ReportData:
"""组装结构化日报数据(M9:写入 MySQL 的前置步骤)。
命名规则: intl_news_daily_{YYYYMMDD_HHMMSS}.html — 支持一天多份日报。
与 news 项目 report_db_design.md 数据契约一致:
- 事件全部归入 `intl` 板块(intl 日报无 xwlb/news/cninfo 板块);
- report_type="intl"、file_name="" → 唯一键退化为 (report_date, intl, ""),
同一天重复生成 = UPDATE 覆盖(幂等);
- 数据总览统计以 JSON 快照存入 stats(前端自行解析)。
"""
events: list[EventRow] = []
for i, ev in enumerate(high_events, 1):
article = ev.get("article", {})
title = (article.get("title_zh") or article.get("title") or "").strip()[:512]
url = article.get("url") or None
# 去重合并后的多来源(如 "Reuters, CNBC"),否则回退单源展示
src = _article_source_label(article)
# 全部来源列表(主源居首,JSON 数组,写入 news_event.sources)
source_list = _article_sources_full(article)
# 归一化:"" / "?" 不入库,留 None(DB 仅存 positive/negative/neutral)
sentiment = ev.get("sentiment") or None
if sentiment in ("", "?"):
sentiment = None
event_type = ev.get("event_type") or None
if event_type in ("", "?"):
event_type = None
events.append(
EventRow(
section="intl",
rank=i,
importance=ev.get("importance") or None,
event_type=event_type,
title=title or "(无标题)",
summary=ev.get("summary_zh") or None,
sentiment=sentiment,
source=src,
sources=source_list,
url=url,
)
)
stats_snapshot: dict[str, Any] = {
"pipeline": {
"raw_total": stats.get("raw_total", 0),
"processed": stats.get("proc", 0),
"deduped": stats.get("deduped", 0),
"embedded": stats.get("emb_count", 0),
"qdrant": stats.get("qdrant_count", 0),
},
"sentiment": {k: v for k, v in sentiments.items() if k not in ("", "?")},
"importance": [
{"importance": k, "count": v} for k, v in sorted(importances.items())
],
"event_types": [
{"event_type": k, "count": v}
for k, v in event_types.most_common(10)
if k not in ("", "?")
],
"source_dist": [
{"source": _source_name(k), "count": v}
for k, v in sources.most_common(15)
if k not in ("", "?")
],
}
return ReportData(
report_date=now.date(),
report_type="intl",
file_name="", # 新生成日报唯一键退化为 (report_date, intl, "")
generated_at=now,
ai_summary=ai_summary or None,
stats=stats_snapshot,
events=events,
)
def generate_report() -> int | None:
"""生成日报并结构化入库(覆盖过去 25 小时数据)。
M9 起完全切换:不再产出 HTML,日报内容写入 MySQL
(news_report / news_event,report_type="intl",file_name="",
同一天重复生成 → 幂等覆盖,不产生多行)。
Returns:
HTML 文件路径,无数据时返回 None
report_id(成功)或 None(无数据/失败)
"""
now = datetime.now()
ts = now.strftime("%Y%m%d_%H%M%S")
date_str = now.strftime("%Y%m%d") # 用于上传目录
logger.info("生成日报: %s(窗口: 过去 %d 小时)", ts, _REPORT_WINDOW_HOURS)
# 加载过去 25 小时数据
@@ -595,13 +717,10 @@ def generate_report() -> Path | None:
)
high = _get_high(all_events, 4)
hi_threshold = 4
if len(high) < 3:
high = _get_high(all_events, 3)
hi_threshold = 3
if len(high) < 3:
high = sorted(all_events, key=lambda e: -e.get("importance", 0))
hi_threshold = 0
# 事件级去重:同一 URL + 同一标题 → 合并
high = _dedup_events(high)
high = high[:_MAX_HIGH_EVENTS]
@@ -609,20 +728,24 @@ def generate_report() -> Path | None:
# AI 摘要
ai_summary = _generate_ai_summary(articles, ts)
# 渲染 HTML
html = _render_html(ts, stats, articles, high, hi_threshold,
sentiments, importances, event_types, sources, ai_summary)
# 结构化入库(替代原 HTML 渲染 + 上传)
report = _build_report_data(
now, stats, high, sentiments, importances, event_types, sources, ai_summary
)
try:
from report_db import connect, save_report
# 本地保存(文件名含时间戳)
_REPORT_DIR.mkdir(parents=True, exist_ok=True)
html_path = _REPORT_DIR / f"intl_news_daily_{ts}.html"
html_path.write_text(html, encoding="utf-8")
logger.info("日报已保存: %s (%d KB)", html_path, len(html) // 1024)
conn = connect()
try:
report_id = save_report(conn, report)
finally:
conn.close()
except Exception as e:
logger.exception("日报入库失败: %s", e)
return None
# 自动上传到日期子目录
_upload_report(html_path, date_str)
return html_path
logger.info("日报已入库: report_id=%s", report_id)
return report_id
# --------------------------------------------------------------------------- #
@@ -644,7 +767,7 @@ def _load_report_config() -> dict:
def _upload_report(html_path: Path, day_str: str) -> bool:
"""上传日报到 Web 服务器(配置来自 system.yaml)。"""
"""上传日报到 Web 服务器(M9 起弃用:日报已改为写库,保留以便回退)。"""
import subprocess
config = _load_report_config()
@@ -678,7 +801,7 @@ def _upload_report(html_path: Path, day_str: str) -> bool:
# --------------------------------------------------------------------------- #
# HTML 渲染
# HTML 渲染(M9 起弃用:日报已改为写库,以下渲染函数保留以便回退)
# --------------------------------------------------------------------------- #
_HTML_TEMPLATE = """<!DOCTYPE html>
@@ -798,7 +921,7 @@ def _render_html(
sources: Counter,
ai_summary: str,
) -> str:
"""组装完整 HTML。"""
"""组装完整 HTML(M9 起弃用,保留以便回退)。"""
# AI 摘要 Markdown → HTML
summary_html = _md_to_html(ai_summary) if ai_summary.strip() else "<p>暂无 AI 摘要</p>"
@@ -915,7 +1038,7 @@ def _dedup_events(events: list[dict]) -> list[dict]:
def _render_event_table(events: list[dict]) -> str:
"""渲染事件表格。"""
"""渲染事件表格(M9 起弃用,保留以便回退)。"""
if not events:
return "<p>暂无符合条件的数据</p>"
+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 各执行一次
# =============================================
# 用法:
# ./scripts/domestic_full.sh # 全新执行
# ./scripts/domestic_full.sh --resume # 从中断处继续(跳过已完成步骤,
# # 步骤状态见 data/run_state/{date}.state)
# =============================================
set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
PROJECT_DIR="$(dirname "$SCRIPT_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 "══════ 国内全流程开始 ══════"
# ── 全流程日志:所有输出实时显示到终端,同时完整写入日志文件 ──
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 ──
export $(grep -v '^#' .env | grep -v '^$' | xargs 2>/dev/null || true)
# ── 1. M1 Pi 抓取(headful Playwright + HTTP 代理)──
LOG "[1/2] Pi M1 抓取(8G headful + HTTP 代理)..."
bash "$SCRIPT_DIR/domestic_crawl_8g.sh" 2>&1 | tail -5 || LOG "WARNING: 部分源抓取失败,继续管道"
# 部分源抓取失败不阻塞管道(原语义);抓取本身按 URL 去重(index.jsonl),
# 已抓取过的 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 管道(含日报)──
LOG "[2/2] 全链路管道..."
bash "$SCRIPT_DIR/pipeline.sh" 2>&1 | tail -10
LOG "━━━ [2/2] M2→M6 管道 + 日报 ━━━"
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 "══════ 国内全流程完成 ✅ ══════"
+71 -20
View File
@@ -2,9 +2,19 @@
# =============================================
# 国内服务器:全链路管道 M2 → M3 → M4 → M5 → M6 → 日报
# =============================================
# 用法:./scripts/pipeline.sh
# 用法:
# ./scripts/pipeline.sh # 全新执行(步骤级不跳过)
# ./scripts/pipeline.sh --resume # 从中断处继续(跳过已完成步骤)
# 前提: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
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
@@ -12,57 +22,98 @@ PROJECT_DIR="$(dirname "$SCRIPT_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
export $(grep -v '^#' .env | grep -v '^$' | xargs 2>/dev/null || true)
# ── M2: 正文提取 ──
echo "[$(date)] [M2] 正文提取..."
.venv/bin/python3 -c "
LOG "━━━ M2 正文提取 ━━━"
step_run M2_extract .venv/bin/python3 -c "
from extractor.pipeline import process_all_sources
stats = process_all_sources()
print(f'M2: {stats[\"total_articles\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s')
"
# ── M3: 去重 ──
echo "[$(date)] [M3] 三层去重..."
.venv/bin/python3 -c "
LOG "━━━ M3 三层去重 ━━━"
step_run M3_dedup .venv/bin/python3 -c "
from dedup.pipeline import dedup_all_sources
stats = dedup_all_sources()
print(f'M3: 唯一 {stats[\"unique\"]}/重复 {stats[\"duplicate\"]}, {stats[\"elapsed_sec\"]:.0f}s')
"
# ── M4: 翻译+事件 ──
echo "[$(date)] [M4] 翻译+事件抽取..."
.venv/bin/python3 -c "
# ── M4: 翻译+事件(AI 大模型)──
LOG "━━━ M4 翻译+事件抽取(AI 大模型: $(ai_model_info translation))━━━"
step_run M4_translate .venv/bin/python3 -c "
from llm.pipeline import translate_all_deduped
stats = translate_all_deduped()
print(f'M4: {stats[\"success\"]}/{stats[\"total\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s')
"
# ── M5: 向量生成 ──
echo "[$(date)] [M5] 向量生成..."
.venv/bin/python3 -c "
# ── M5: 向量生成(AI 大模型)──
LOG "━━━ M5 向量生成(AI 大模型: $(ai_model_info embedding))━━━"
step_run M5_embed .venv/bin/python3 -c "
from embedding.pipeline import embed_all_events
stats = embed_all_events()
print(f'M5: {stats[\"success\"]}/{stats[\"total\"]} 篇, {stats[\"elapsed_sec\"]:.0f}s')
"
# ── M6: Qdrant 入库 ──
echo "[$(date)] [M6] Qdrant 入库..."
.venv/bin/python3 -c "
LOG "━━━ M6 Qdrant 入库 ━━━"
step_run M6_index .venv/bin/python3 -c "
from vectorstore.pipeline import ingest_all_embeddings
stats = ingest_all_embeddings()
print(f'M6: {stats[\"ingested\"]}/{stats[\"total\"]} 条, {stats[\"elapsed_sec\"]:.0f}s')
"
# ── 日报 ──
echo "[$(date)] [日报] 生成日报..."
.venv/bin/python3 -c "
# ── 日报(AI 大模型)──
LOG "━━━ 日报生成(AI 大模型: $(ai_model_info daily_report))━━━"
step_run report .venv/bin/python3 -c "
from scheduler.reporter import generate_report
path = generate_report()
print(f'日报: {path}')
report_id = generate_report()
print(f'日报: report_id={report_id}' if report_id is not None else '日报: 无数据/失败')
"
echo "[$(date)] ═══ 全链路管道完成 ✅ ═══"
LOG "══════ 全链路管道完成 ✅ ═══════"
+99
View File
@@ -1,5 +1,6 @@
"""M3 三层去重模块单元测试。"""
import json
from pathlib import Path
import pytest
@@ -17,6 +18,7 @@ from dedup import (
normalize_content,
simhash64,
)
from dedup.pipeline import dedup_source
from extractor.models import ProcessedArticle
# --------------------------------------------------------------------------- #
@@ -531,3 +533,100 @@ class TestDedupResult:
summary = r.short_summary()
assert "[DUP/simhash]" 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.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):
monkeypatch.delenv("DASHSCOPE_API_KEY", raising=False)
monkeypatch.setenv("QWEN_API_KEY", "sk-qwen-key")
+29
View File
@@ -457,6 +457,35 @@ class TestLoadLLMConfig:
with pytest.raises(ValueError, match="API key 未配置"):
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)
+216
View File
@@ -0,0 +1,216 @@
"""日报结构化入库:模型 + _build_report_data 组装(纯逻辑,不连 DB)。"""
from __future__ import annotations
import json
from collections import Counter
from datetime import date, datetime
import pytest
from pydantic import ValidationError
from report_db.models import EventRow, ReportData
from scheduler.reporter import _build_report_data
# --------------------------------------------------------------------------- #
# 模型默认值
# --------------------------------------------------------------------------- #
class TestEventRow:
def test_minimal(self) -> None:
ev = EventRow(section="intl", rank=1, title="标题")
assert ev.importance is None
assert ev.sentiment is None
def test_full(self) -> None:
ev = EventRow(
section="intl", rank=2, importance=4, event_type="地缘政治",
title="t", summary="s", sentiment="negative", source="InvestingLive",
url="https://x.com/1",
)
assert ev.sentiment == "negative"
def test_missing_title_raises(self) -> None:
with pytest.raises(ValidationError):
EventRow(section="intl", rank=1) # type: ignore[call-arg]
class TestReportData:
def test_defaults(self) -> None:
r = ReportData(
report_date=date(2026, 8, 4),
report_type="intl",
generated_at=datetime(2026, 8, 4, 8, 0),
)
assert r.file_name == ""
assert r.stats == {}
assert r.events == []
def test_with_events(self) -> None:
r = ReportData(
report_date=date(2026, 8, 4),
report_type="intl",
generated_at=datetime(2026, 8, 4, 8, 0),
ai_summary="摘要",
events=[EventRow(section="intl", rank=1, title="t")],
)
assert len(r.events) == 1
# --------------------------------------------------------------------------- #
# _build_report_data(intl 结构)
# --------------------------------------------------------------------------- #
def _fake_high_event(title_zh: str, importance: int, *,
event_type: str = "其他", sentiment: str = "neutral",
summary_zh: str = "摘要", source_id: str = "investinglive",
url: str = "https://investinglive.com/news/1") -> dict:
"""构造 _dedup_events 之后的事件结构(含 article 属性)。"""
return {
"importance": importance,
"sentiment": sentiment,
"summary_zh": summary_zh,
"event_type": event_type,
"stock_codes": [],
"article": {
"title": f"EN {title_zh}",
"title_zh": title_zh,
"url": url,
"source_id": source_id,
},
}
class TestBuildReportData:
def test_sections_and_ranks(self) -> None:
now = datetime(2026, 8, 4, 8, 0, 0)
high = [
_fake_high_event("新闻A", 5),
_fake_high_event("新闻B", 4),
_fake_high_event("新闻C", 3),
]
r = _build_report_data(
now, {"raw_total": 100, "proc": 90}, high,
Counter(), Counter(), Counter(), Counter(), "AI摘要",
)
assert r.report_date == date(2026, 8, 4)
assert r.report_type == "intl"
assert r.file_name == ""
assert r.ai_summary == "AI摘要"
assert [(e.section, e.rank) for e in r.events] == [
("intl", 1), ("intl", 2), ("intl", 3),
]
def test_field_mapping(self) -> None:
now = datetime(2026, 8, 4, 8, 0, 0)
high = [_fake_high_event("英伟达财报超预期", 5, event_type="公司财报",
sentiment="positive", summary_zh="营收超预期",
source_id="investinglive",
url="https://investinglive.com/news/42")]
r = _build_report_data(now, {}, high, Counter(), Counter(), Counter(),
Counter(), "")
ev = r.events[0]
assert ev.title == "英伟达财报超预期" # title_zh 优先
assert ev.importance == 5
assert ev.event_type == "公司财报"
assert ev.summary == "营收超预期"
assert ev.sentiment == "positive"
assert ev.url == "https://investinglive.com/news/42"
assert ev.source == "InvestingLive" # 域名 → 源展示名
def test_source_fallback_without_url(self) -> None:
now = datetime(2026, 8, 4, 8, 0, 0)
ev = _fake_high_event("无URL事件", 4, url="")
r = _build_report_data(now, {}, [ev], Counter(), Counter(), Counter(),
Counter(), "")
assert r.events[0].url is None
# url 为空时回退 source_id → sources.yaml 展示名
assert r.events[0].source == "InvestingLive"
def test_stats_snapshot(self) -> None:
now = datetime(2026, 8, 4, 8, 0, 0)
high = [
_fake_high_event("A", 5, event_type="宏观", sentiment="positive"),
_fake_high_event("B", 4, event_type="宏观", sentiment="negative"),
]
r = _build_report_data(
now,
{"raw_total": 120, "proc": 100, "deduped": 90,
"emb_count": 80, "qdrant_count": 75},
high,
Counter({"positive": 1, "negative": 1}),
Counter({5: 1, 4: 1}),
Counter({"宏观": 2}),
Counter({"investinglive": 2}),
"摘要",
)
# stats 可 JSON 序列化(入库时 json.dumps)
json.dumps(r.stats, ensure_ascii=False)
assert r.stats["pipeline"]["raw_total"] == 120
assert r.stats["sentiment"] == {"positive": 1, "negative": 1}
assert r.stats["importance"] == [
{"importance": 4, "count": 1}, {"importance": 5, "count": 1},
]
assert r.stats["event_types"] == [{"event_type": "宏观", "count": 2}]
assert r.stats["source_dist"] == [{"source": "InvestingLive", "count": 2}]
def test_title_truncated(self) -> None:
now = datetime(2026, 8, 4, 8, 0, 0)
ev = _fake_high_event("长" * 600, 4)
r = _build_report_data(now, {}, [ev], Counter(), Counter(), Counter(),
Counter(), "")
assert len(r.events[0].title) == 512
def test_empty_events(self) -> None:
now = datetime(2026, 8, 4, 8, 0, 0)
r = _build_report_data(now, {}, [], Counter(), Counter(), Counter(),
Counter(), "")
assert r.events == []
assert r.ai_summary is None
def test_question_mark_sentiment_normalized(self) -> None:
now = datetime(2026, 8, 4, 8, 0, 0)
ev = _fake_high_event("未分类情绪", 4, sentiment="?")
r = _build_report_data(now, {}, [ev], Counter(), Counter(), Counter(),
Counter(), "")
# "?" 不写入 DB,留 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"]