# 数据管道设计 ## 概述 从 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) | ``(新增+修改)+ ``(删除) | **唯一每日数据源** | | 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 覆盖) → → 标记 retracted=True ``` ### 放弃 E-utilities 搜索的原因 | 维度 | 旧: E-utilities API 多路搜索 | 新: FTP 每日更新文件 | |------|------------------------------|---------------------| | 完整性 | 仅返回匹配搜索词的结果。宽搜词覆盖不到的文章→**永恒漏掉** | **100%** — 当天所有增/改/删的记录都在一个文件里 | | 删除检测 | 搜不到已删的记录 → 无法标记删除 | `` 明确列出被删 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 位顺序号,每年重置) 文件包含三种顶层元素: - `` — **新增**或**已修改**的文献记录 - `` — 图书章节(跳过,非期刊文献) - `` — 被删除的 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 — `` 独立处理 | | 是否适用于 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.954.248"] 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" ``` **三层过滤逻辑**(`_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 时的文本兜底)。 --- ## 二、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 解析 │ │ │ │ 对每条 提取(复用 _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}] │ │ │ │ 对每条 : │ │ 查本地 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()) │ │ │ │ 对每篇文章按三层逐一判断: │ │ │ │ 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 → ❌ 丢弃 │ │ │ │ 过滤结果:日均 ~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. 解析 → _extract_article() → _is_oncology() 过滤 4. 解析 → 标记 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 篇/req,3 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 不再需要 |