# 数据管道设计 ## 概述 从 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 篇/req,3 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 页写入缓存,避免分页碎片