2026-07-27 07:59:18 +08:00
|
|
|
|
# 数据管道设计
|
|
|
|
|
|
|
|
|
|
|
|
## 概述
|
|
|
|
|
|
|
|
|
|
|
|
从 PubMed FTP 每日增量拉取文献 → 按肿瘤科 MeSH 过滤 → MeSH 映射标签引擎打标 → 入库 → 用户匹配 → 推送 Feed。全过程自动化,每天凌晨执行。
|
|
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
|
|
|
|
|
|
## 一、PubMed 数据获取
|
|
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
### 1.1 运行模式(决策:2026-07-25)
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
数据源统一为 **FTP 每日更新文件**(`ftp://ftp.ncbi.nlm.nih.gov/pubmed/updatefiles/pubmed25updateNNN.xml.gz`),取代原有的 E-utilities API 多路搜索策略。
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
| 数据源 | 协议 | 内容 | 用途 |
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|--------|------|------|------|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
| **FTP 每日更新文件** | FTP (gzip XML) | `<PubmedArticle>`(新增+修改)+ `<DeleteCitation>`(删除) | **唯一每日数据源** |
|
|
|
|
|
|
| NCBI E-utilities | REST XML | 搜索 + 抓取 | **仅降级后备**(FTP 不可用时) |
|
|
|
|
|
|
| Europe PMC | REST JSON | 搜索 + 抓取 | 仅基线导入后一次性回填(`rebuild_from_baseline.py`) |
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
```
|
2026-07-27 08:35:12 +08:00
|
|
|
|
每日更新(03:07 UTC):
|
|
|
|
|
|
FTP 下载 pubmed25updateNNN.xml.gz → lxml 解析
|
|
|
|
|
|
→ _extract_article() 提取字段(复用 baseline 解析器)
|
|
|
|
|
|
→ _is_oncology() 肿瘤科过滤(复用 baseline 过滤逻辑)
|
|
|
|
|
|
→ _update_lit_from_article() upsert(幂等,同 PMID 覆盖)
|
|
|
|
|
|
→ <DeleteCitation> → 标记 retracted=True
|
2026-07-27 07:59:18 +08:00
|
|
|
|
```
|
|
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
### 放弃 E-utilities 搜索的原因
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
| 维度 | 旧: E-utilities API 多路搜索 | 新: FTP 每日更新文件 |
|
|
|
|
|
|
|------|------------------------------|---------------------|
|
|
|
|
|
|
| 完整性 | 仅返回匹配搜索词的结果。宽搜词覆盖不到的文章→**永恒漏掉** | **100%** — 当天所有增/改/删的记录都在一个文件里 |
|
|
|
|
|
|
| 删除检测 | 搜不到已删的记录 → 无法标记删除 | `<DeleteCitation>` 明确列出被删 PMID |
|
|
|
|
|
|
| MeSH 过渡 | 需 retagger 回查 in-process→medline 变化 | 更新文件直接包含 MeSH 状态变化后的完整记录 |
|
|
|
|
|
|
| 限速 | 3 req/s(免费) → 大结果集分页慢 | 单文件下载,无 API 限速 |
|
|
|
|
|
|
| 代码复杂度 | 2 套搜索(MAJR + 宽搜)+ 3 条解析路径 + retagger | 1 套下载 + 1 套解析(复用 baseline) |
|
|
|
|
|
|
| 定时任务 | 精搜 + 宽搜 + retagger = 3 个 cron | 1 个 cron |
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
FTP baseline 导入(`scripts/pubmed_baseline.py`)用于**首次全量数据**,完成后每日增量数据完全由 FTP update files 接管。
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
### 1.2 FTP 每日更新文件结构
|
|
|
|
|
|
|
|
|
|
|
|
文件名格式:`pubmed25updateNNNN.xml.gz`(NNNN 为 4 位顺序号,每年重置)
|
|
|
|
|
|
|
|
|
|
|
|
文件包含三种顶层元素:
|
|
|
|
|
|
- `<PubmedArticle>` — **新增**或**已修改**的文献记录
|
|
|
|
|
|
- `<PubmedBookArticle>` — 图书章节(跳过,非期刊文献)
|
|
|
|
|
|
- `<DeleteCitation>` — 被删除的 PMID(含 PubMed/PMC 两种来源)
|
|
|
|
|
|
|
|
|
|
|
|
NLM 每天发布一个新文件(约 5-20MB 压缩后),全年约 365 个文件。
|
|
|
|
|
|
|
|
|
|
|
|
### 1.3 基线过渡(Baseline → FTP 衔接)
|
|
|
|
|
|
|
|
|
|
|
|
基线导入和 FTP 每日更新之间存在**年份衔接问题**:
|
|
|
|
|
|
|
|
|
|
|
|
| 问题 | 说明 | 处理方式 |
|
|
|
|
|
|
|------|------|---------|
|
|
|
|
|
|
| 基线年份 | 基线文件 `pubmed26n*.xml.gz` 中的 "26" 表示 2026 年 | **自动检测**:`rebuild_from_baseline.py` 导入完成后自动记录 `pipeline_runs(run_type='baseline_import', run_metadata={'year': 2026})` |
|
|
|
|
|
|
| FTP 文件前缀 | 2026 年的增量文件为 `pubmed26updateNNNN.xml.gz` | 从基线记录读取年份,构造前缀 `pubmed{year}update` |
|
|
|
|
|
|
| 首次运行 | 基线刚导入完,FTP 文件可能从年初就开始有了 | 按 `sequence > 0` 拉取所有文件(与基线无交叠问题) |
|
|
|
|
|
|
| 跨年份过渡 | 12 月→1 月,文件前缀从 `pubmed25update` → `pubmed26update` | 前一年没文件时自动尝试上年;从 `pipeline_runs` 的 `run_metadata.last_sequence` 判断跨年 |
|
|
|
|
|
|
|
|
|
|
|
|
具体流程:
|
|
|
|
|
|
```
|
|
|
|
|
|
首次运行(无 checkpoint):
|
|
|
|
|
|
1. pipeline_runs 查 baseline_import 记录 → year=2026
|
|
|
|
|
|
2. 构造前缀 pubmed26update → 列出 FTP 文件
|
|
|
|
|
|
3. 全部下载处理(sequence 从 0001 开始)
|
|
|
|
|
|
4. 记录 run_metadata = {ftp_year: 2026, last_sequence: 最大序号}
|
|
|
|
|
|
|
|
|
|
|
|
后续运行(有 checkpoint):
|
|
|
|
|
|
1. 读上次 run_metadata → ftp_year=2026, last_sequence=1234
|
|
|
|
|
|
2. 列出 pubmed26update*,筛选 sequence > 1234
|
|
|
|
|
|
3. 处理 → 更新 run_metadata.last_sequence
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
### 1.4 检查点(Checkpoint)机制
|
|
|
|
|
|
|
|
|
|
|
|
使用 `pipeline_runs` 表的多字段组合追踪已处理状态:
|
|
|
|
|
|
|
|
|
|
|
|
| 字段 | 用途 |
|
|
|
|
|
|
|------|------|
|
|
|
|
|
|
| `run_type` = `"daily_ftp_update"` | 标记 FTP 更新运行 |
|
|
|
|
|
|
| `processed_date` (DATE) | 最后一个已处理的 EDAT 日期 |
|
|
|
|
|
|
| `run_metadata.ftp_year` (JSON) | FTP 文件年份(决定文件前缀) |
|
|
|
|
|
|
| `run_metadata.last_sequence` (JSON) | 最后一个已处理的文件序号 |
|
|
|
|
|
|
| `articles_new` | 新增文献数 |
|
|
|
|
|
|
| `articles_updated` | 更新文献数 |
|
|
|
|
|
|
| `articles_deleted` | DeleteCitation 处理数 |
|
|
|
|
|
|
|
|
|
|
|
|
**检查点行为:**
|
|
|
|
|
|
- 每天运行时,读 `pipeline_runs WHERE run_type='daily_ftp_update' ORDER BY created_at DESC LIMIT 1`
|
|
|
|
|
|
- 从上次 `last_sequence` 的序号+1开始拉取
|
|
|
|
|
|
- 如果当天无新 update 文件发布,跳过(非错误)
|
|
|
|
|
|
- 运行失败 → 不写 PipelineRun → 下次重跑同序列号(幂等)
|
|
|
|
|
|
- FTP 文件含多天的累积更新 → 按 `_update_lit_from_article()` 覆盖不会重复
|
|
|
|
|
|
|
|
|
|
|
|
### 1.5 EDAT 作为判断依据
|
|
|
|
|
|
|
|
|
|
|
|
EDAT(Entrez Date)是 PubMed 记录首次进入 Entrez 系统的日期,是判断"新记录"的最可靠依据:
|
|
|
|
|
|
|
|
|
|
|
|
| 特性 | EDAT |
|
|
|
|
|
|
|------|------|
|
|
|
|
|
|
| 是否可变 | **不可变** — 记录一旦入库,EDAT 终生不变 |
|
|
|
|
|
|
| 精度 | 精确到日(Y/M/D),无时间分量 → DATE 类型完全匹配 |
|
|
|
|
|
|
| 唯一性 | 同一天两条记录可能有不同 EDAT?不可能 — 所有新记录取当天日期 |
|
|
|
|
|
|
| 删除检测 | 不依赖 EDAT — `<DeleteCitation>` 独立处理 |
|
|
|
|
|
|
| 是否适用于 FTP | ✅ FTP 文件中可提取每条记录的 History/PubMedPubDate[@PubStatus="entrez"] |
|
|
|
|
|
|
|
|
|
|
|
|
EDAT 对应字段:`GlobalLiterature.entrez_date`(DATE,新增字段,迁移:`1421ea169bb6`)。
|
|
|
|
|
|
|
|
|
|
|
|
### 1.5 配置文件驱动过滤
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
```yaml
|
|
|
|
|
|
# config/specialties/oncology.yaml
|
|
|
|
|
|
pubmed_filter:
|
|
|
|
|
|
mesh_include_categories: ["C04"]
|
|
|
|
|
|
mesh_include_subcategories: ["C04.557", "C04.588", "C04.697"]
|
2026-07-28 11:02:29 +08:00
|
|
|
|
mesh_cross_include: ["E02.319", "E02.815", "D27.505.954.248"]
|
2026-07-27 07:59:18 +08:00
|
|
|
|
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"
|
|
|
|
|
|
```
|
|
|
|
|
|
|
2026-07-28 11:02:29 +08:00
|
|
|
|
**三层过滤逻辑**(`_is_oncology()` 实现):
|
|
|
|
|
|
|
|
|
|
|
|
| 层级 | 条件 | 动作 |
|
|
|
|
|
|
|------|------|------|
|
|
|
|
|
|
| **Step 0** `pub_type_exclude` | PublicationType ∈ {Editorial, Letter, News, Comment} | ❌ 直接丢弃 |
|
|
|
|
|
|
| **Step 1** MeSH 树号匹配 | MeSH heading 在 C04(Neoplasms)或 cross-include 树号(E02.319 抗癌药、E02.815 放疗、D27.505.954.248 抗肿瘤药)下 | ✅ 收录 |
|
|
|
|
|
|
| **Step 2** 文本评分兜底 | 无 MeSH 标引的文献:标题含关键词 → 收录;摘要含 ≥2 个关键词 → 收录;摘要单关键词 → 放过 | ✅/❌ 按评分 |
|
|
|
|
|
|
|
|
|
|
|
|
三层按顺序执行:Step 0 最先(直接排除),Step 1 中间(MeSH 精确匹配),Step 2 最后(无 MeSH 时的文本兜底)。
|
|
|
|
|
|
|
2026-07-27 07:59:18 +08:00
|
|
|
|
---
|
|
|
|
|
|
|
|
|
|
|
|
## 二、Pipeline 完整流程
|
|
|
|
|
|
|
|
|
|
|
|
```
|
2026-07-27 08:35:12 +08:00
|
|
|
|
ARQ 定时任务 (每日 FTP 更新 03:07 UTC)
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│
|
|
|
|
|
|
▼
|
|
|
|
|
|
┌─────────────────────────────────────────────────────────┐
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ Step 1: FTP 文件下载 + 检查点检测 │
|
|
|
|
|
|
│ 从 pipeline_runs 读上次处理的 processed_date │
|
|
|
|
|
|
│ 下载 pubmed25updateNNNN.xml.gz(序号 > 上次处理的) │
|
|
|
|
|
|
│ 失败 → 降级到 E-utilities reldate=1&datetype=edat 后备 │
|
|
|
|
|
|
│ 预计耗时:10-30 秒(gzip ~10MB 解压后 ~200MB) │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
└───────────────────┬─────────────────────────────────────┘
|
|
|
|
|
|
│
|
|
|
|
|
|
▼
|
|
|
|
|
|
┌─────────────────────────────────────────────────────────┐
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ Step 2: XML 解析 │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ 对每条 <PubmedArticle> 提取(复用 _extract_article()): │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ - PMID (PubMed ID) │
|
|
|
|
|
|
│ - Title, Abstract │
|
|
|
|
|
|
│ - Authors [{given, family, affiliation}] │
|
|
|
|
|
|
│ - DOI, Journal (name, ISSN, ISOAbbreviation, │
|
|
|
|
|
|
│ volume, issue, pages) │
|
|
|
|
|
|
│ - PublicationDate (year, month, day) │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ - History (entrez_date, create_date, pubmed_revised,│
|
|
|
|
|
|
│ date_completed) │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ - PublicationType (期刊文章/综述/RCT/指南...) │
|
|
|
|
|
|
│ - MeSH Headings [{descriptor, ui, major}] │
|
|
|
|
|
|
│ - Language │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ - ChemicalList [{registry_number, ui, name}] │
|
|
|
|
|
|
│ - GeneSymbolList [string] │
|
|
|
|
|
|
│ - NumberOfReferences (int) │
|
|
|
|
|
|
│ - PublicationStatus (epublish/ppublish/aheadofprint)│
|
|
|
|
|
|
│ - ArticleDate (电子出版日期, DateType="Electronic") │
|
|
|
|
|
|
│ - DataBankList [{name, accession_numbers[]}] │
|
|
|
|
|
|
│ - SupplMeshList [{ui, name, type}] │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ 对每条 <DeleteCitation>: │
|
|
|
|
|
|
│ 查本地 PMID → 标记 retracted=True + updated_at │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ 预计耗时:30-60 秒(单个文件 ~500-5000 条记录) │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ 单一解析器(统一维护): │
|
|
|
|
|
|
│ - _extract_article() — FTP, lxml (pubmed_baseline.py)│
|
|
|
|
|
|
│ - 注:pubmed_api.py 中的_parse_pubmed_xml() │
|
|
|
|
|
|
│ 是 E-utilities 降级路径,需与 _extract_article() │
|
|
|
|
|
|
│ 保持同步 │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ 入库路径(已存在): │
|
|
|
|
|
|
│ - _update_lit_from_article() — upsert 更新已有记录 │
|
|
|
|
|
|
│ - 新 PMID 通过 GlobalLiterature() 构造函数创建 │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ │
|
|
|
|
|
|
└───────────────────┬─────────────────────────────────────┘
|
|
|
|
|
|
│
|
|
|
|
|
|
▼
|
|
|
|
|
|
┌─────────────────────────────────────────────────────────┐
|
2026-07-28 11:02:29 +08:00
|
|
|
|
│ Step 3: 肿瘤科过滤(三层,配置驱动,复用 _is_oncology()) │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ │
|
2026-07-28 11:02:29 +08:00
|
|
|
|
│ 对每篇文章按三层逐一判断: │
|
|
|
|
|
|
│ │
|
|
|
|
|
|
│ Step 0 — 排除类型检查: │
|
|
|
|
|
|
│ if publication_type ∈ {Editorial, Letter, News, │
|
|
|
|
|
|
│ Comment} │
|
|
|
|
|
|
│ → ❌ 丢弃(无需往下判断) │
|
|
|
|
|
|
│ │
|
|
|
|
|
|
│ Step 1 — MeSH 树号匹配: │
|
|
|
|
|
|
│ elif 任一 heading 在 C04(Neoplasms)下 → ✅ 收录 │
|
|
|
|
|
|
│ elif 任一 heading 在 cross-include 树号下 │
|
|
|
|
|
|
│ (E02.319 抗癌药 / E02.815 放疗 /│
|
|
|
|
|
|
│ D27.505.954.248 抗肿瘤药) → ✅ 收录 │
|
|
|
|
|
|
│ │
|
|
|
|
|
|
│ Step 2 — 文本评分兜底(无 MeSH 标引时): │
|
|
|
|
|
|
│ elif 标题含关键词 → ✅ 收录 │
|
|
|
|
|
|
│ elif 摘要含 ≥2 个关键词 → ✅ 收录 │
|
|
|
|
|
|
│ else → ❌ 丢弃 │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ 过滤结果:日均 ~5000 条总记录 → 约 400-800 条肿瘤相关 │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ 预计耗时: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, │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
│ tag_id, is_major) ON CONFLICT DO UPDATE │
|
|
|
|
|
|
│ │
|
|
|
|
|
|
│ DeleteCitation 处理: │
|
|
|
|
|
|
│ UPDATE global_literature SET retracted=TRUE, │
|
|
|
|
|
|
│ updated_at=NOW() WHERE pmid IN (...) │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
│ │
|
|
|
|
|
|
│ 预计耗时:< 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 分钟(千级用户,百级文献/用户) │
|
2026-07-27 08:35:12 +08:00
|
|
|
|
└───────────────────────┬─────────────────────────────────┘
|
|
|
|
|
|
│
|
|
|
|
|
|
▼
|
|
|
|
|
|
┌─────────────────────────────────────────────────────────┐
|
|
|
|
|
|
│ Step 7: 更新检查点 │
|
|
|
|
|
|
│ INSERT INTO pipeline_runs (run_type='daily_ftp_update',│
|
|
|
|
|
|
│ processed_date=今天, status='success', │
|
|
|
|
|
|
│ articles_new=..., articles_updated=..., │
|
|
|
|
|
|
│ articles_deleted=...) │
|
2026-07-27 07:59:18 +08:00
|
|
|
|
└─────────────────────────────────────────────────────────┘
|
|
|
|
|
|
|
|
|
|
|
|
总计耗时:每日常规 5-10 分钟
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
|
|
|
|
|
|
## 三、性能估算
|
|
|
|
|
|
|
|
|
|
|
|
| 环节 | 数据量 | 预计耗时 |
|
|
|
|
|
|
|------|--------|:---:|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
| FTP 文件下载 + 解压 | 1 个文件,~10MB gzip | 10-30 秒 |
|
|
|
|
|
|
| XML 解析(含 DeleteCitation) | 单文件 ~500-5000 条 | 30-60 秒 |
|
2026-07-27 07:59:18 +08:00
|
|
|
|
| MeSH 过滤 | — | 1-2 分钟 |
|
|
|
|
|
|
| 打标 | 每篇 5-15 个 MeSH | 1-2 分钟 |
|
2026-07-27 08:35:12 +08:00
|
|
|
|
| 入库 + 引用查询 | 400-800 ON CONFLICT + elink | 1-2 分钟 |
|
|
|
|
|
|
| DeleteCitation 更新 | 通常 < 10 条/天 | < 1 秒 |
|
2026-07-27 07:59:18 +08:00
|
|
|
|
| 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 = [
|
2026-07-27 08:35:12 +08:00
|
|
|
|
"app.tasks.pubmed_pipeline.daily_ftp_update", # FTP 每日增量(取代精搜+宽搜+retagger)
|
|
|
|
|
|
"app.tasks.pubmed_pipeline.daily_citation_update", # 每日引用刷新
|
|
|
|
|
|
"app.tasks.pubmed_pipeline.daily_digest_task", # 每日摘要邮件推送
|
2026-07-27 07:59:18 +08:00
|
|
|
|
"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 = [
|
2026-07-27 08:35:12 +08:00
|
|
|
|
# 每日 FTP 增量更新(取代 daily_pubmed_pipeline + weekly_broad_pipeline)
|
|
|
|
|
|
cron(daily_ftp_update, minute=7, hour=3),
|
2026-07-27 07:59:18 +08:00
|
|
|
|
# 每日引用刷新
|
|
|
|
|
|
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
|
2026-07-27 08:35:12 +08:00
|
|
|
|
async def daily_ftp_update(ctx):
|
|
|
|
|
|
"""每日 FTP 增量更新(取代旧 daily_pubmed_pipeline + weekly_broad_pipeline)
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
流程:
|
|
|
|
|
|
1. 查 pipeline_runs 获取上次处理的 EDAT 检查点
|
|
|
|
|
|
2. 下载 FTP pubmed25update* 文件(序号 > 上次处理的值)
|
|
|
|
|
|
3. 解析 <PubmedArticle> → _extract_article() → _is_oncology() 过滤
|
|
|
|
|
|
4. 解析 <DeleteCitation> → 标记 retracted=True
|
|
|
|
|
|
5. upsert → _update_lit_from_article() 或 GlobalLiterature() 新建
|
|
|
|
|
|
6. MeSH 映射 → tag_article() 打标
|
|
|
|
|
|
7. 生成 Feed → generate_feeds_for_literature()
|
|
|
|
|
|
8. 更新检查点 → INSERT pipeline_runs
|
|
|
|
|
|
"""
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
async def daily_citation_update(ctx):
|
|
|
|
|
|
"""刷新最近文献被引次数"""
|
|
|
|
|
|
# NCBI elink 接口,200 篇/req,3 req/s
|
|
|
|
|
|
...
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
|
|
|
|
|
|
## 五、错误处理与恢复
|
|
|
|
|
|
|
|
|
|
|
|
| 场景 | 处理 |
|
|
|
|
|
|
|------|------|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
| FTP 连接超时 | 重试 3 次,间隔 60s,仍失败则降级到 E-utilities reldate |
|
|
|
|
|
|
| FTP 文件破损(gzip 校验失败) | 跳过该文件,记录日志,下次重试 |
|
|
|
|
|
|
| XML 解析错误 | 跳过单条 record,记录错误 PMID |
|
2026-07-27 07:59:18 +08:00
|
|
|
|
| 入库冲突 | `ON CONFLICT (pmid) DO UPDATE`,幂等 |
|
|
|
|
|
|
| MeSH 映射找不到 | 记录 `unmapped_mesh_terms` 日志,定期人工审查 |
|
|
|
|
|
|
| 用户匹配超时 | 分批处理,每批 500 用户 |
|
2026-07-27 08:35:12 +08:00
|
|
|
|
| 整管道失败 | `processed_date` 不更新 → 次日重跑同一天数据(幂等保证) |
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
## 六、代码实现
|
|
|
|
|
|
|
|
|
|
|
|
### 6.1 文件清单
|
|
|
|
|
|
|
|
|
|
|
|
| 文件 | 状态 | 说明 |
|
|
|
|
|
|
|------|------|------|
|
|
|
|
|
|
| `app/services/pubmed_daily_update.py` | ✅ **新增** | 核心服务:FTP 下载、解析、过滤、upsert、删除标记、检查点 |
|
|
|
|
|
|
| `app/tasks/worker.py` | ✅ **已改** | cron 替换:`daily_ftp_update` 取代 `daily_pubmed_pipeline` + `weekly_broad_pipeline` + `daily_mesh_retagger` |
|
|
|
|
|
|
| `app/api/v1/admin.py` | ✅ **已改** | `POST /pipeline/run` 默认 mode 改为 `ftp` |
|
|
|
|
|
|
| `app/models/operations.py` | ✅ **已改** | `PipelineRun.processed_date` DATE + `run_metadata` JSON(检查点) |
|
|
|
|
|
|
| `alembic/versions/f3b135d62407_...` | ✅ **已迁移** | 新增 `pipeline_runs.processed_date` 列 |
|
|
|
|
|
|
| `alembic/versions/0ee585329fc6_...` | ✅ **已迁移** | 新增 `pipeline_runs.metadata` JSON 列 |
|
|
|
|
|
|
|
|
|
|
|
|
### 6.2 核心流程
|
|
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
run_daily_ftp_update()
|
|
|
|
|
|
│
|
|
|
|
|
|
├─ 1. get_last_checkpoint()
|
|
|
|
|
|
│ SELECT run_metadata FROM pipeline_runs
|
|
|
|
|
|
│ WHERE run_type='daily_ftp_update' AND status='success'
|
|
|
|
|
|
│ ORDER BY created_at DESC LIMIT 1
|
|
|
|
|
|
│ → {ftp_year, last_sequence} 或 None(首次运行)
|
|
|
|
|
|
│
|
|
|
|
|
|
├─ 2. 确定文件年份:checkpoint.ftp_year ?? baseline_year ?? 当前年
|
|
|
|
|
|
│
|
|
|
|
|
|
├─ 3. list_update_files(year, after_sequence)
|
|
|
|
|
|
│ FTP NLST → pubmed{year}updateNNNN.xml.gz 列表
|
|
|
|
|
|
│ 筛选 sequence > last_sequence
|
|
|
|
|
|
│
|
|
|
|
|
|
├─ 4. 逐个 download_file() + gzip.decompress()
|
|
|
|
|
|
│
|
|
|
|
|
|
├─ 5. process_update_file()
|
|
|
|
|
|
│ ├─ DeleteCitation PMIDs → UPDATE retracted=TRUE
|
|
|
|
|
|
│ ├─ _extract_article() 解析(复用 pubmed_baseline.py)
|
|
|
|
|
|
│ ├─ _is_oncology() 过滤(复用 pubmed_baseline.py)
|
|
|
|
|
|
│ ├─ upsert: 已存在 → _update_lit_from_article() + tag_article()
|
|
|
|
|
|
│ └─ 新建 → GlobalLiterature() + tag_article() + generate_feeds()
|
|
|
|
|
|
│
|
|
|
|
|
|
└─ 6. PipelineRun(processed_date=当天, run_metadata={ftp_year, last_sequence})
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
### 6.3 降级策略
|
|
|
|
|
|
|
|
|
|
|
|
FTP 不可用时(网络故障、NLM FTP 宕机),回退到 E-utilities `run_eutils_fallback()`:
|
|
|
|
|
|
```python
|
|
|
|
|
|
async def run_eutils_fallback(last_checkpoint):
|
|
|
|
|
|
# 使用 reldate=1&datetype=edat 查询最近 1 天新增
|
|
|
|
|
|
# 调用 _run_pipeline(use_majr=True, use_broad=True)
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
|
|
|
|
|
|
## 七、数据质量监控
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
```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';
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
## 八、热搜缓存管道
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
### 8.1 概述
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
首页侧栏"热搜"展示近期高被引文献(支持分页 + "换一批")。首页文献列表也使用预缓存策略。两者都避免每次请求走全套查询。
|
|
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
### 8.2 数据流
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
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 → 写回缓存 → 返回
|
|
|
|
|
|
```
|
|
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
### 8.3 关键文件
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
| 文件 | 说明 |
|
|
|
|
|
|
|------|------|
|
|
|
|
|
|
| `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 优先,不可用时降级为内存缓存 |
|
|
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
### 8.4 配置参数
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
| 参数 | 值 | 说明 |
|
|
|
|
|
|
|------|-----|------|
|
|
|
|
|
|
| `HOT_LIMIT` | 30 | 缓存保留的文献总数(最多 30 条) |
|
|
|
|
|
|
| `CACHE_TTL` | 1800s (30min) | 缓存有效期 |
|
|
|
|
|
|
| `page_size` | 5 | 每页展示数 |
|
|
|
|
|
|
| 最大页数 | 3(前端限制,后端支持更多) |
|
|
|
|
|
|
|
2026-07-27 08:35:12 +08:00
|
|
|
|
### 8.5 缓存失效说明
|
2026-07-27 07:59:18 +08:00
|
|
|
|
|
|
|
|
|
|
- 定时任务每 30 分钟覆盖写入,完全刷新缓存
|
|
|
|
|
|
- 缓存通过 `CACHE_TTL` 自动过期,Redis 重启后首次请求回退实时查询
|
|
|
|
|
|
- 实时查询回退时仅第 1 页写入缓存,避免分页碎片
|
2026-07-27 08:35:12 +08:00
|
|
|
|
|
|
|
|
|
|
---
|
|
|
|
|
|
|
|
|
|
|
|
## 九、MeSH 年更维护
|
|
|
|
|
|
|
|
|
|
|
|
### 9.1 背景
|
|
|
|
|
|
|
|
|
|
|
|
NLM 每年 11 月中–12 月中进行 MeSH 词汇年更(新词加入、旧词合并、树号变更)。期间 indexed 记录暂停灌入,只发 in-process/publisher 记录。年更过后,大量已有记录的 MeSH 标引会被刷新。
|
|
|
|
|
|
|
|
|
|
|
|
### 9.2 刷新判断
|
|
|
|
|
|
|
|
|
|
|
|
`meshed_date` 字段记录 **MeSH 数据的版本时间**(NLM 实际完成 MeSH 标引的时间),不是我们的刷新操作时间:
|
|
|
|
|
|
|
|
|
|
|
|
| 状态 | meshed_date | 含义 |
|
|
|
|
|
|
|------|-------------|------|
|
|
|
|
|
|
| 已标引 | = date_completed | 当前 MeSH 是最新的 |
|
|
|
|
|
|
| 年更后已刷 | = 新的 date_completed | 已通过年更刷新 |
|
|
|
|
|
|
| 待刷新 | < 年更新阈值 | MeSH 可能已过时 |
|
|
|
|
|
|
| 无 MeSH | NULL | in-process/publisher,无标引 |
|
|
|
|
|
|
|
|
|
|
|
|
> ⚠️ `meshed_date` 始终使用 NLM 端的时间戳(`date_completed`),不写 `datetime.now()`。仅在 NLM 未返回 `date_completed` 的极少数情况下降级使用当前时间。
|
|
|
|
|
|
|
|
|
|
|
|
### 9.3 年更期间的应对
|
|
|
|
|
|
|
|
|
|
|
|
FTP 每日更新文件在年更期间的行为:
|
|
|
|
|
|
- Indexed 记录**暂停**出现在更新文件中(NLM 暂停新标引)
|
|
|
|
|
|
- In-process 和 publisher 记录**照常**出现在更新文件中
|
|
|
|
|
|
- 年更结束后,重标引的记录**一次性在更新文件中出现**(当天的 update 文件会异常大)
|
|
|
|
|
|
|
|
|
|
|
|
这意味着 **FTP 方案天然适配年更** — 不需要特殊的管道逻辑,更新文件里有什么就处理什么。
|
|
|
|
|
|
|
|
|
|
|
|
### 9.4 刷新操作(手工执行)
|
|
|
|
|
|
|
|
|
|
|
|
> **⚠️ 此操作为手工任务**。不要自动执行。需先观察 PubMed 状态判断年更是否完成。
|
|
|
|
|
|
|
|
|
|
|
|
判断时机:
|
|
|
|
|
|
1. 观察 NLM 公告(https://www.nlm.nih.gov/pubs/techbull/)
|
|
|
|
|
|
2. 检查 pipeline 日志:indexed 记录是否恢复灌入
|
|
|
|
|
|
3. 确认年更窗口结束(通常 12 月中)
|
|
|
|
|
|
|
|
|
|
|
|
执行命令(admin token 需要有 `platform_operator` 角色):
|
|
|
|
|
|
|
|
|
|
|
|
```bash
|
|
|
|
|
|
# 查询待刷新文献数
|
|
|
|
|
|
psql -U scilit_dev -d scilit_dev -c "
|
|
|
|
|
|
SELECT count(*) FROM global_literature
|
|
|
|
|
|
WHERE meshed_date IS NOT NULL
|
|
|
|
|
|
AND meshed_date < '2025-11-01'; # 去年年更前
|
|
|
|
|
|
|
|
|
|
|
|
# 全量刷新 MeSH(按 PMID 批次 re-fetch)
|
|
|
|
|
|
curl -X POST localhost:8000/api/v1/admin/pipeline/run \
|
|
|
|
|
|
-H "Authorization: Bearer $(admin_token)"
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
> Pipeline 正常运行时 `_update_lit_from_article()` 会自动覆盖 `meshed_date = date_completed`。但对于年更场景,重抓后 MeSH 已变,NLM 返回的 date_completed 仍是原值,所以 pipeline 更新不会正确标记。正确的年更刷新机制为一个独立脚本(待实现):按 `meshed_date < 年更新阈值` + `mesh_headings IS NOT NULL` 条件批量 efetch → 对比 MeSH 是否有变 → 有变则更新并用 `now()` 覆盖 `meshed_date`。
|
|
|
|
|
|
|
|
|
|
|
|
### 9.5 待实现脚本(低优先级)
|
|
|
|
|
|
|
|
|
|
|
|
| 任务 | 文件 | 优先级 | 说明 |
|
|
|
|
|
|
|------|------|--------|------|
|
|
|
|
|
|
| MeSH 年更批量刷新脚本 | `scripts/refresh_mesh_annual.py` | P2 | 按 meshed_date 阈值批量 efetch,检测 MeSH 变化,更新 meshed_date=now() |
|
|
|
|
|
|
| — | `app/services/mesh_retagger.py` | ❌ **已废弃** | FTP 每日更新文件直接包含 in-process→medline 过渡的完整记录,retagger 不再需要 |
|