Files
backend/docs/06-数据管道设计.md
T
34047007@qq.com 62ca8fa6b8
CI / backend (push) Canceled after 0s
CI / frontend (push) Canceled after 0s
chore: batch commit remaining changes
Includes search engine improvements, Alembic migrations,
new services (pubmed_daily_update, query_expansion),
frontend updates, and documentation sync.
2026-07-27 08:35:12 +08:00

585 lines
30 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 运行模式(决策:2026-07-25
数据源统一为 **FTP 每日更新文件**`ftp://ftp.ncbi.nlm.nih.gov/pubmed/updatefiles/pubmed25updateNNN.xml.gz`),取代原有的 E-utilities API 多路搜索策略。
| 数据源 | 协议 | 内容 | 用途 |
|--------|------|------|------|
| **FTP 每日更新文件** | FTP (gzip XML) | `<PubmedArticle>`(新增+修改)+ `<DeleteCitation>`(删除) | **唯一每日数据源** |
| NCBI E-utilities | REST XML | 搜索 + 抓取 | **仅降级后备**FTP 不可用时) |
| Europe PMC | REST JSON | 搜索 + 抓取 | 仅基线导入后一次性回填(`rebuild_from_baseline.py` |
```
每日更新(03:07 UTC):
FTP 下载 pubmed25updateNNN.xml.gz → lxml 解析
→ _extract_article() 提取字段(复用 baseline 解析器)
→ _is_oncology() 肿瘤科过滤(复用 baseline 过滤逻辑)
→ _update_lit_from_article() upsert(幂等,同 PMID 覆盖)
→ <DeleteCitation> → 标记 retracted=True
```
### 放弃 E-utilities 搜索的原因
| 维度 | 旧: 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 |
FTP baseline 导入(`scripts/pubmed_baseline.py`)用于**首次全量数据**,完成后每日增量数据完全由 FTP update files 接管。
### 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 作为判断依据
EDATEntrez 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 配置文件驱动过滤
```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 定时任务 (每日 FTP 更新 03:07 UTC)
┌─────────────────────────────────────────────────────────┐
│ Step 1: FTP 文件下载 + 检查点检测 │
│ 从 pipeline_runs 读上次处理的 processed_date │
│ 下载 pubmed25updateNNNN.xml.gz(序号 > 上次处理的) │
│ 失败 → 降级到 E-utilities reldate=1&datetype=edat 后备 │
│ 预计耗时:10-30 秒(gzip ~10MB 解压后 ~200MB
└───────────────────┬─────────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 2: XML 解析 │
│ │
│ 对每条 <PubmedArticle> 提取(复用 _extract_article()): │
│ - PMID (PubMed ID) │
│ - Title, Abstract │
│ - Authors [{given, family, affiliation}] │
│ - DOI, Journal (name, ISSN, ISOAbbreviation, │
│ volume, issue, pages) │
│ - PublicationDate (year, month, day) │
│ - History (entrez_date, create_date, pubmed_revised,│
│ date_completed) │
│ - 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}] │
│ │
│ 对每条 <DeleteCitation>
│ 查本地 PMID → 标记 retracted=True + updated_at │
│ │
│ 预计耗时:30-60 秒(单个文件 ~500-5000 条记录) │
│ │
│ 单一解析器(统一维护): │
│ - _extract_article() — FTP, lxml (pubmed_baseline.py)│
│ - 注:pubmed_api.py 中的_parse_pubmed_xml() │
│ 是 E-utilities 降级路径,需与 _extract_article() │
│ 保持同步 │
│ │
│ 入库路径(已存在): │
│ - _update_lit_from_article() — upsert 更新已有记录 │
│ - 新 PMID 通过 GlobalLiterature() 构造函数创建 │
│ │
└───────────────────┬─────────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 3: 肿瘤科过滤(配置驱动,复用 _is_oncology()
│ │
│ 对每篇文章的 MeSH Headings[] 逐一检查: │
│ if 任一 heading.ui 在 mesh_include_categories 中 │
│ → 保留 │
│ elif 任一 heading.ui 在 mesh_cross_include 中 │
│ → 保留(交叉领域) │
│ else │
│ → 丢弃 │
│ │
│ 过滤结果:日均 ~5000 条总记录 → 约 400-800 条肿瘤相关 │
│ 预计耗时: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) ON CONFLICT DO UPDATE │
│ │
│ DeleteCitation 处理: │
│ UPDATE global_literature SET retracted=TRUE, │
│ updated_at=NOW() WHERE pmid IN (...) │
│ │
│ 预计耗时:< 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 分钟(千级用户,百级文献/用户) │
└───────────────────────┬─────────────────────────────────┘
┌─────────────────────────────────────────────────────────┐
│ Step 7: 更新检查点 │
│ INSERT INTO pipeline_runs (run_type='daily_ftp_update',│
│ processed_date=今天, status='success', │
│ articles_new=..., articles_updated=..., │
│ articles_deleted=...) │
└─────────────────────────────────────────────────────────┘
总计耗时:每日常规 5-10 分钟
```
---
## 三、性能估算
| 环节 | 数据量 | 预计耗时 |
|------|--------|:---:|
| FTP 文件下载 + 解压 | 1 个文件,~10MB gzip | 10-30 秒 |
| XML 解析(含 DeleteCitation | 单文件 ~500-5000 条 | 30-60 秒 |
| MeSH 过滤 | — | 1-2 分钟 |
| 打标 | 每篇 5-15 个 MeSH | 1-2 分钟 |
| 入库 + 引用查询 | 400-800 ON CONFLICT + elink | 1-2 分钟 |
| DeleteCitation 更新 | 通常 < 10 条/天 | < 1 秒 |
| 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_ftp_update", # FTP 每日增量(取代精搜+宽搜+retagger)
"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 = [
# 每日 FTP 增量更新(取代 daily_pubmed_pipeline + weekly_broad_pipeline
cron(daily_ftp_update, minute=7, hour=3),
# 每日引用刷新
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_ftp_update(ctx):
"""每日 FTP 增量更新(取代旧 daily_pubmed_pipeline + weekly_broad_pipeline
流程:
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
"""
async def daily_citation_update(ctx):
"""刷新最近文献被引次数"""
# NCBI elink 接口,200 篇/req3 req/s
...
```
---
## 五、错误处理与恢复
| 场景 | 处理 |
|------|------|
| FTP 连接超时 | 重试 3 次,间隔 60s,仍失败则降级到 E-utilities reldate |
| FTP 文件破损(gzip 校验失败) | 跳过该文件,记录日志,下次重试 |
| XML 解析错误 | 跳过单条 record,记录错误 PMID |
| 入库冲突 | `ON CONFLICT (pmid) DO UPDATE`,幂等 |
| MeSH 映射找不到 | 记录 `unmapped_mesh_terms` 日志,定期人工审查 |
| 用户匹配超时 | 分批处理,每批 500 用户 |
| 整管道失败 | `processed_date` 不更新 → 次日重跑同一天数据(幂等保证) |
---
## 六、代码实现
### 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)
```
---
## 七、数据质量监控
```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';
```
---
## 八、热搜缓存管道
### 8.1 概述
首页侧栏"热搜"展示近期高被引文献(支持分页 + "换一批")。首页文献列表也使用预缓存策略。两者都避免每次请求走全套查询。
### 8.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 → 写回缓存 → 返回
```
### 8.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 优先,不可用时降级为内存缓存 |
### 8.4 配置参数
| 参数 | 值 | 说明 |
|------|-----|------|
| `HOT_LIMIT` | 30 | 缓存保留的文献总数(最多 30 条) |
| `CACHE_TTL` | 1800s (30min) | 缓存有效期 |
| `page_size` | 5 | 每页展示数 |
| 最大页数 | 3(前端限制,后端支持更多) |
### 8.5 缓存失效说明
- 定时任务每 30 分钟覆盖写入,完全刷新缓存
- 缓存通过 `CACHE_TTL` 自动过期,Redis 重启后首次请求回退实时查询
- 实时查询回退时仅第 1 页写入缓存,避免分页碎片
---
## 九、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 不再需要 |