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

19 KiB
Raw Blame History

数据管道设计

概述

从 PubMed FTP 每日增量拉取文献 → 按肿瘤科 MeSH 过滤 → MeSH 映射标签引擎打标 → 入库 → 用户匹配 → 推送 Feed。全过程自动化,每天凌晨执行。


一、PubMed 数据获取

1.1 实际运行模式:API 优先

当前部署使用 NCBI E-utilities APIEurope 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 配置文件驱动过滤

# 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 任务实现

# 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 用户
整管道失败 告警通知 → 次日重试前一天的数据

六、数据质量监控

-- 查看每天的采集量(异常检测)
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 页写入缓存,避免分页碎片