Files
backend/docs/06-数据管道设计.md
T
34047007@qq.com a6cd99a4ca
CI / backend (push) Canceled after 0s
CI / frontend (push) Canceled after 0s
feat: initial commit - oncology literature search platform
OncoLit: a multi-tenant oncology literature search, feed, and
collaboration platform. Built with FastAPI + Vue 3 + PostgreSQL.
Includes PubMed pipeline, drug approvals, AI summaries, and
systematic review tools.
2026-07-27 07:59:18 +08:00

385 lines
19 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 数据管道设计
## 概述
从 PubMed FTP 每日增量拉取文献 → 按肿瘤科 MeSH 过滤 → MeSH 映射标签引擎打标 → 入库 → 用户匹配 → 推送 Feed。全过程自动化,每天凌晨执行。
---
## 一、PubMed 数据获取
### 1.1 实际运行模式:API 优先
当前部署使用 **NCBI E-utilities API****Europe PMC API** 作为数据源(`app/services/pubmed_api.py`),而非 FTP baseline。
| 数据源 | 协议 | 限速 | 用途 |
|--------|------|------|------|
| NCBI E-utilities | REST XML | 3 req/s(免费)/ 10 req/s(有 API Key | 每日精搜(MAJR |
| Europe PMC | REST JSON | 无严格限速,支持游标分页 | MAJR 精搜 + 标题/摘要宽搜 |
```
精搜(每日 03:07 UTC):
ONCOLOGY_SEARCH_QUERIES × [MAJR] → NCBI esearch → efetch
→ 约 20 篇/query × 10 queries = 200 篇/天
宽搜(每周日 03:37 UTC):
BROAD_ONCOLOGY_QUERIES × Title/Abstract → Europe PMC 游标分页
→ 覆盖 in-process + publisher 记录
→ 与精搜通过 seen_pmids 自动去重
```
### 1.2 暂未使用 FTP
FTP baseline 导入(`scripts/pubmed_baseline.py`)已完成开发但尚未在生产中执行,原因:
| 维度 | API | FTP |
|------|-----|-----|
| 初始数据量 | ~200 篇/天自动积累 | 一次导入 ~30GB |
| 部署复杂度 | 零额外配置 | 需下载 1200+ 个 XML 文件 |
| 推荐时机 | 每日自动运行 | 仅在需要全量历史文献时手动触发 |
### 1.3 配置文件驱动过滤
```yaml
# config/specialties/oncology.yaml
pubmed_filter:
mesh_include_categories: ["C04"]
mesh_include_subcategories: ["C04.557", "C04.588", "C04.697"]
mesh_cross_include: ["E02.319", "E02.815", "D27.505"]
pub_type_prefer:
- "Guideline"
- "Randomized Controlled Trial"
- "Meta-Analysis"
- "Systematic Review"
- "Clinical Trial, Phase III"
- "Clinical Trial, Phase II"
- "Review"
- "Case Reports"
pub_type_exclude:
- "Editorial"
- "Letter"
- "News"
- "Comment"
```
---
## 二、Pipeline 完整流程
```
ARQ 定时任务 (每日精搜 03:07 UTC / 每周宽搜 周日 03:37 UTC)
┌─────────────────────────────────────────────────────────┐
│ Step 1: API 搜索 │
│ NCBI E-utilities (esearch + efetch) — MAJR 精搜 │
│ Europe PMC REST API — 标题/摘要宽搜 │
│ 精搜 ~200 篇/天,宽搜 ~2000 篇/周日 │
│ 自动去重 (seen_pmids 集合) │
│ 预计耗时:1-2 分钟 │
└───────────────────┬─────────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 2: XML/JSON 解析 │
│ │
│ 对每条 PubMedArticle 提取: │
│ - PMID (PubMed ID) │
│ - Title, Abstract │
│ - Authors [{given, family, affiliation}] │
│ - DOI, Journal (name, ISSN, ISOAbbreviation, │
│ volume, issue, pages) │
│ - PublicationDate (year, month, day) │
│ - PublicationType (期刊文章/综述/RCT/指南...) │
│ - MeSH Headings [{descriptor, ui, major}] │
│ - Language │
│ - ChemicalList [{registry_number, ui, name}] │ ← 增强字段
│ - GeneSymbolList [string] │ ← 增强字段
│ - NumberOfReferences (int) │ ← 增强字段
│ - PublicationStatus (epublish/ppublish/aheadofprint)│ ← 增强字段
│ - ArticleDate (电子出版日期, DateType="Electronic") │ ← 增强字段
│ - DataBankList [{name, accession_numbers[]}] │ ← 增强字段
│ - SupplMeshList [{ui, name, type}] │ ← 增强字段
│ │
│ 增强字段覆盖率(实测 1369 篇): │
│ publication_status ~100% article_date ~73% │
│ chemical_list ~37% num_refs ~91% │
│ gene_symbols ~10% databank_list ~25% │
│ suppl_mesh_list ~35% │
│ │
│ 特殊处理:DeleteCitation → 标记删除不处理 │
│ 预计耗时:1-2 分钟 (2000-4000条) │
│ │
│ 三处解析器(需保持同步): │
│ - _parse_pubmed_xml() — NCBI E-utilities, ElementTree│
│ - _extract_article() — FTP baseline, lxml │
│ - _parse_europe_pmc_article() — Europe PMC, JSON │
│ │
│ 两条入库路径(均需传递增强字段): │
│ - _process_article() — 创建新 GlobalLiterature() │
│ - _update_lit_from_article() — 更新已有记录 │
│ │
└───────────────────┬─────────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 3: 肿瘤科过滤(配置驱动) │
│ │
│ 对每篇文章的 MeSH Headings[] 逐一检查: │
│ if 任一 heading.ui 在 mesh_include_categories 中 │
│ → 保留 │
│ elif 任一 heading.ui 在 mesh_cross_include 中 │
│ → 保留(交叉领域) │
│ else │
│ → 丢弃 │
│ │
│ 过滤结果:日均 3000 条 → 约 400-600 条肿瘤相关 │
│ 预计耗时:1-2 分钟 │
└───────────────────┬─────────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 4: MeSH → Tag 映射 + 自动打标 │
│ │
│ 对保留的每篇文章: │
│ 对每个 MeSH heading: │
│ a. 查 global_tags WHERE mesh_ui = heading.ui │
│ b. 如果匹配 → 生成 (tag_id, is_major) │
│ c. 同时提取 tag_category 用于后续分组 │
│ │
│ 额外补充标签(Phase 2+): │
│ - 从 PubType 映射研究类型标签 │
│ - 从 journal.issn 映射期刊等级标签 │
│ - Phase 3: LLM 辅助判断癌种亚型/靶点 │
│ │
│ 预计耗时:1-2 分钟(每篇 5-15 个 MeSH
└───────────────────┬─────────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 5: 入库 PostgreSQL │
│ │
│ 对每条过滤后的文章: │
│ INSERT INTO global_literature ... ON CONFLICT (pmid) │
│ DO UPDATE (如果已存在则更新) │
│ INSERT INTO global_literature_tags (literature_id, │
│ tag_id, is_major) │
│ │
│ 预计耗时:< 1 分钟(批量 INSERT
└───────────────────┬─────────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 5b: PMC OA 全文抓取(可选) │
│ │
│ 对新入库且有 PMC ID 的文献,自动抓取 OA 全文: │
│ elink 查 PMID → PMCID 映射 │
│ efetch 获取 JATS XML │
│ jats_parser 解析 → sections / tables / refs / chems │
│ 存入 full_text_sections (JSON) + license_info │
│ 同时预抽取基线特征表 (Table 1 → pre_extracted_data) │
│ │
│ 覆盖:~15% 的 PubMed 文献有 PMC 全文 │
│ 限速:同 NCBI E-utilities 池 (3/10 req/s) │
│ 预计耗时:每篇 1-3 秒 │
└───────────────────┬─────────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 5c: 期刊去重同步(_ensure_journal
│ │
│ 对每条入库的文献,自动确保 global_journals 中有对应记录: │
│ 1. 按 journal_issn 精确匹配 → 命中则跳过 │
│ 2. 未命中 → 写新行,ON CONFLICT (issn) DO NOTHING │
│ 3. 无 ISSN 的 → name 精确匹配 + pg_trgm 模糊兜底 │
│ │
│ ISOAbbreviation (journal_iso) 同步写入 global_literature│
│ 预计耗时:< 1 秒(批量 upsert
└───────────────────┬─────────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 6: 用户匹配引擎 │
│ │
│ 获取所有活跃用户 → 查 user_subscriptions │
│ 对每条新文献: │
│ article_tags = 文献已打的标签集合 │
│ user_tags = 用户关注的标签集合 │
│ │
│ 匹配策略: │
│ loose: article_tags ∩ user_tags ≠ ∅ │
│ standard: article_major_tags ∩ user_tags ≠ ∅ │
│ strict: user_tags ⊆ article_tags │
│ │
│ 命中 → INSERT user_feed (user_id, literature_id, │
│ matched_tags, priority) │
│ │
│ priority 计算: │
│ - 全命中+顶刊 → must_read │
│ - 部分命中 → recommended │
│ - 领域匹配 → related │
│ │
│ 预计耗时:< 1 分钟(千级用户,百级文献/用户) │
└─────────────────────────────────────────────────────────┘
总计耗时:每日常规 5-10 分钟
```
---
## 三、性能估算
| 环节 | 数据量 | 预计耗时 |
|------|--------|:---:|
| API 搜索 + 抓取 | 200-2000 篇/次 | 1-2 分钟 |
| XML/JSON 解析 | 2000-4000 篇/天 | 1-2 分钟 |
| MeSH 过滤 | — | 1-2 分钟 |
| 打标 | 每篇 5-15 个 MeSH | 1-2 分钟 |
| 入库 + 引用查询 | 400-600 INSERT/UPDATE + elink | 1-2 分钟 |
| OA 全文抓取 | ~15% 有 PMC | 每篇 1-3 秒 |
| 期刊去重同步 | 每篇 | < 1 秒 |
| 用户匹配 | 千级用户 | < 1 分钟 |
| **总计** | | **5-10 分钟** |
---
## 四、ARQ 任务实现
```python
# app/tasks/worker.py
from arq import cron
from arq.connections import RedisSettings
class WorkerSettings:
functions = [
"app.tasks.pubmed_pipeline.daily_pubmed_pipeline",
"app.tasks.pubmed_pipeline.weekly_broad_pipeline",
"app.tasks.pubmed_pipeline.daily_citation_update",
"app.tasks.pubmed_pipeline.daily_digest_task",
"app.tasks.pubmed_pipeline.refresh_hot_articles_cache",
"app.tasks.pubmed_pipeline.refresh_homepage_feed",
]
redis_settings = RedisSettings(
host=settings.REDIS_HOST,
port=settings.REDIS_PORT,
password=settings.REDIS_PASSWORD,
)
cron_jobs = [
# 每日精搜(MAJR MeSH 高精度,20 篇/query
cron(daily_pubmed_pipeline, minute=7, hour=3),
# 每周宽搜(Title/Abstract 覆盖 in-process + publisher
cron(weekly_broad_pipeline, minute=37, hour=3, day_of_week=0),
# 每日引用刷新
cron(daily_citation_update, minute=13, hour=5),
# 每日摘要邮件推送
cron(daily_digest_task, minute=30, hour=22),
# 热搜缓存刷新(每 30 分钟)
cron(refresh_hot_articles_cache, minute="*/30"),
# 首页 Feed 缓存刷新(每 30 分钟)
cron(refresh_homepage_feed, minute="*/30"),
]
# app/tasks/pubmed_pipeline.py
async def daily_pubmed_pipeline(ctx):
"""每日 PubMed 精搜"""
# 通过 pubmed_api.py _run_pipeline(use_majr=True, use_broad=False) 执行
# 约 10 queries × 20 篇 = 200 篇/天
...
async def weekly_broad_pipeline(ctx):
"""每周宽搜"""
# 通过 pubmed_api.py _run_pipeline(use_majr=False, use_broad=True) 执行
# 覆盖 in-process + publisher 记录,自动 seen_pmids 去重
...
async def daily_citation_update(ctx):
"""刷新最近文献被引次数"""
# NCBI elink 接口,200 篇/req3 req/s
...
```
---
## 五、错误处理与恢复
| 场景 | 处理 |
|------|------|
| FTP 连接超时 | 重试 3 次,间隔 60s,仍失败则告警 |
| XML 解析错误 | 跳过单条,记录错误 PMID |
| 入库冲突 | `ON CONFLICT (pmid) DO UPDATE`,幂等 |
| MeSH 映射找不到 | 记录 `unmapped_mesh_terms` 日志,定期人工审查 |
| 用户匹配超时 | 分批处理,每批 500 用户 |
| 整管道失败 | 告警通知 → 次日重试前一天的数据 |
---
## 六、数据质量监控
```sql
-- 查看每天的采集量(异常检测)
SELECT
date(created_at) as day,
count(*) as articles,
count(DISTINCT journal) as journals
FROM global_literature
WHERE created_at > now() - interval '30 days'
GROUP BY day ORDER BY day DESC;
-- 检查未匹配 MeSH 的文献
SELECT pmid, title FROM global_literature
WHERE mesh_headings = '[]' OR mesh_headings IS NULL
AND created_at > now() - interval '7 days';
```
---
## 七、热搜缓存管道
### 7.1 概述
首页侧栏"热搜"展示近期高被引文献(支持分页 + "换一批")。首页文献列表也使用预缓存策略。两者都避免每次请求走全套查询。
### 7.2 数据流
```
ARQ 定时任务(每 30 分钟):
┌─ refresh_hot_articles_cache()
│ 查 DB(近 1 年有被引的文献,按 cited_by_count DESC 取前 30 条)
│ 写入 Redis (key: "hot_articles", TTL: 1800s)
└─ refresh_homepage_feed()
查 DB(最新 50 条,不取 JSON 大字段,仅 ~10KB/条)
写入 Redis (key: "homepage:feed", TTL: 1800s)
API 请求:
GET /public/hot-articles?page=1&page_size=5 → 优先读 Redis
GET /public/feed?page=1 → 优先读 Redis
未命中 → 回退实时查询 DB → 写回缓存 → 返回
```
### 7.3 关键文件
| 文件 | 说明 |
|------|------|
| `app/services/hot_articles_cache.py` | `precompute_hot_articles()` — 预计算写入缓存;`_build_cards()` — 数据格式转换 |
| `app/api/v1/public.py` | `GET /hot-articles` — 先读缓存,未命中回退实时查询 |
| `app/tasks/worker.py` | `refresh_hot_articles_cache` — 定时任务(`*/30` 分钟 UTC |
| `app/core/cache.py` | `CacheService` — Redis 优先,不可用时降级为内存缓存 |
### 7.4 配置参数
| 参数 | 值 | 说明 |
|------|-----|------|
| `HOT_LIMIT` | 30 | 缓存保留的文献总数(最多 30 条) |
| `CACHE_TTL` | 1800s (30min) | 缓存有效期 |
| `page_size` | 5 | 每页展示数 |
| 最大页数 | 3(前端限制,后端支持更多) |
### 7.5 缓存失效说明
- 定时任务每 30 分钟覆盖写入,完全刷新缓存
- 缓存通过 `CACHE_TTL` 自动过期,Redis 重启后首次请求回退实时查询
- 实时查询回退时仅第 1 页写入缓存,避免分页碎片