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 正常入库
This commit is contained in:
@@ -140,11 +140,27 @@ 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`),经翻译透传后在日报事件 `source` 展示(如 "Barron's, CNBC, Reuters",最多 3 个)。
|
||||||
|
|
||||||
### 日报入库(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`。
|
||||||
|
|||||||
+41
-3
@@ -40,18 +40,56 @@ 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-v4-flash"
|
deepseek_model: "deepseek-v4-flash"
|
||||||
qwen_model: "qwen-plus"
|
qwen_model: "qwen-plus"
|
||||||
timeout_sec: 60
|
timeout_sec: 60
|
||||||
max_attempts: 3 # 单篇总尝试次数(含首次),失败后指数退避重试;日报 AI 摘要同用此值
|
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"
|
||||||
|
|||||||
+34
-1
@@ -1,6 +1,39 @@
|
|||||||
# continuation.md — English Financial News 项目状态
|
# continuation.md — English Financial News 项目状态
|
||||||
|
|
||||||
> 最后更新:2026-08-04
|
> 最后更新:2026-08-12
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 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
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
|||||||
@@ -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:
|
||||||
|
|||||||
+6
-1
@@ -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)",
|
||||||
|
)
|
||||||
|
|||||||
+22
-8
@@ -29,14 +29,26 @@ _QWEN_DEFAULT_MODEL = "qwen-plus"
|
|||||||
_DEFAULT_MAX_ATTEMPTS = 3
|
_DEFAULT_MAX_ATTEMPTS = 3
|
||||||
|
|
||||||
|
|
||||||
def _load_system_config() -> dict:
|
def _load_system_config(scene: str | None = None) -> dict:
|
||||||
"""加载 configs/system.yaml 中 llm 段配置。"""
|
"""加载 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 {}
|
||||||
@@ -64,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")
|
||||||
|
|||||||
@@ -264,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,
|
||||||
@@ -332,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,
|
||||||
|
|||||||
@@ -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 = "" # 英文原标题
|
||||||
|
|||||||
+2
-2
@@ -119,8 +119,8 @@ 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)
|
||||||
template = PromptTemplate()
|
template = PromptTemplate()
|
||||||
|
|
||||||
|
|||||||
+22
-4
@@ -96,6 +96,23 @@ def _url_source_label(url: str, source_id: str = "") -> str:
|
|||||||
return domain_map.get(domain, domain)
|
return domain_map.get(domain, domain)
|
||||||
|
|
||||||
|
|
||||||
|
def _article_source_label(article: dict) -> str | None:
|
||||||
|
"""文章来源展示:去重合并后多来源时拼接展示名,否则回退单源逻辑。
|
||||||
|
|
||||||
|
多来源(ProcessedArticle.source_ids 长度 > 1)时展示如 "Reuters, CNBC",
|
||||||
|
最多取前 3 个来源,截断至 64 字符(news_event.source 为 VARCHAR(64))。
|
||||||
|
"""
|
||||||
|
src_ids = list(dict.fromkeys(
|
||||||
|
s for s in (article.get("source_ids") or []) if s and s != "?"
|
||||||
|
))
|
||||||
|
if len(src_ids) > 1:
|
||||||
|
names = list(dict.fromkeys(
|
||||||
|
(_source_name(s) or s) for s in src_ids[:3]
|
||||||
|
))
|
||||||
|
return ", ".join(names)[:64]
|
||||||
|
return _url_source_label(article.get("url"), article.get("source_id", ""))
|
||||||
|
|
||||||
|
|
||||||
# 日报覆盖时间窗口(小时)
|
# 日报覆盖时间窗口(小时)
|
||||||
_REPORT_WINDOW_HOURS = 25
|
_REPORT_WINDOW_HOURS = 25
|
||||||
|
|
||||||
@@ -299,11 +316,11 @@ def _call_llm_simple(
|
|||||||
"""
|
"""
|
||||||
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"])
|
||||||
|
|
||||||
config = _llm_client_cache["config"]
|
config = _llm_client_cache["config"]
|
||||||
@@ -321,7 +338,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()
|
||||||
@@ -581,7 +598,8 @@ 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)
|
||||||
# 归一化:"" / "?" 不入库,留 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 ("", "?"):
|
||||||
|
|||||||
@@ -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"]
|
||||||
|
|||||||
@@ -463,6 +463,29 @@ class TestLoadLLMConfig:
|
|||||||
config = load_llm_config(provider="deepseek")
|
config = load_llm_config(provider="deepseek")
|
||||||
assert config.max_attempts == 3
|
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)
|
||||||
|
|||||||
@@ -177,3 +177,25 @@ 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)。"""
|
||||||
|
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 含两个来源
|
||||||
|
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].url == "https://reuters.com/news/9"
|
||||||
|
|
||||||
|
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"
|
||||||
|
|||||||
Reference in New Issue
Block a user