fix: 生产高级搜索 500 — 预算超时降级 + pub_date 排序索引 + affiliations_text trgm 反范式列
CI / backend (push) Waiting to run
CI / frontend (push) Waiting to run

- count/主查询独立短预算,超时降级不再 500:count 超时返回 _COUNT_CAP(10 万+)
  而非 0(修复"找到 0 条结果"、分页被隐藏),主查询超时降级空结果
- 年份筛选改用 pub_date 范围,主查询 ORDER BY pub_date DESC NULLS LAST, id DESC
  走新索引 ix_gl_pub_date_sort(Index Scan,替代 Parallel Seq Scan + Sort,89.5 万行 40s→<0.2s)
- affiliations_text 反规范化列 + gin_trgm 索引 + 触发器维护,替换 7 处 jsonb lateral
  EXISTS(病理查询慢 ~40x);存量由 scripts/backfill_affiliations.py 分批后台回填
- ATM 门控 _has_structured_terms:结构化查询(含字段标签)的 plain 词不做 MeSH 展开,
  修复 "lung neoplasms[MAJR]+年份" count 5s + main 6s 双超时(11.72s→2.02s)
- query_expansion/mesh IN(subquery) 去显式 DISTINCT(百万级 id 集合物化排序)
- alembic env 改 connect()(transaction_per_migration 每迁移独立事务,支持迁移内分批提交)
This commit is contained in:
34047007@qq.com
2026-08-10 21:48:42 +08:00
parent ab64d2aa0a
commit a4d27da8ea
9 changed files with 516 additions and 89 deletions
+3 -1
View File
@@ -40,7 +40,9 @@ async def run_async_migrations() -> None:
configuration = config.get_section(config.config_ini_section, {})
configuration["sqlalchemy.url"] = settings.DATABASE_URL
connectable = async_engine_from_config(configuration, prefix="sqlalchemy.", poolclass=pool.NullPool)
async with connectable.begin() as conn:
# 用 connect()(而非 begin()):不隐式开外层事务,让 alembic 按 transaction_per_migration
# 每迁移独立事务;也才能让迁移内的 autocommit_block()(分批回填批间提交)按预期工作。
async with connectable.connect() as conn:
await conn.run_sync(do_run_migrations)
await connectable.dispose()
@@ -0,0 +1,30 @@
"""add ix_gl_pub_date_sort for search order-by
搜索主查询排序索引:ORDER BY pub_date DESC NULLS LAST, id DESCkeyset 分页 + LIMIT 21)。
ASC 的 ix_gl_pub_date_covering 无法服务 DESC NULLS LAST 顺序(方向与 NULL 序均不匹配),
导致主查询退化为 Parallel Seq Scan + Sort89.5 万行 ~40s,触发生产 500)。
此索引让 LIMITed 主查询走 Index Scan + Filter,常规查询 <0.2s。
Revision ID: 235096e73c2a
Revises: 55105f0bb1d7
Create Date: 2026-08-10 19:17:11.237529
"""
from typing import Sequence, Union
from alembic import op
revision: str = '235096e73c2a'
down_revision: Union[str, None] = '55105f0bb1d7'
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
op.execute("""
CREATE INDEX IF NOT EXISTS ix_gl_pub_date_sort
ON global_literature (pub_date DESC NULLS LAST, id DESC)
""")
def downgrade() -> None:
op.execute("DROP INDEX IF EXISTS ix_gl_pub_date_sort")
@@ -0,0 +1,141 @@
"""add affiliations_text column for trgm-indexed affiliation search
新增 affiliations_text 列(从 authors JSONB 提取 affiliation 值拼接),
由 update_literature_search_tsv() 触发器维护。替代搜索的 jsonb_array_elements 逐行
lateral EXISTS(病理查询 40x 成本),让 affiliation ILIKE 走 pg_trgm 索引。
与 author_names_textf1a2b3c4d5e6)同模式。触发器函数只追加 affiliations_text 赋值,
保留当前 search_tsv 全部成分(title/abstract/author_names/chemical/gene/mesh/keywords)。
⚠️ 回填不在此迁移内:全表 UPDATE 会持 AccessExclusiveLock 整表(本地 578 万行 >12min
阻塞搜索读)。列先空建,回填由 scripts/backfill_affiliations.py 分批后台执行
(每批 2 万行、批间提交释放锁,幂等 WHERE affiliations_text IS NULL)。
Revision ID: 3355d1f9dbba
Revises: 235096e73c2a
Create Date: 2026-08-10 20:11:00.000000
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
revision: str = '3355d1f9dbba'
down_revision: Union[str, None] = '235096e73c2a'
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
# 当前 update_literature_search_tsv() 完整定义(从库中 pg_get_functiondef 采集),
# 仅追加 NEW.affiliations_text 赋值,其余(含 mesh/keywords 成分)原样保留。
_FUNC_NEW = """
CREATE OR REPLACE FUNCTION update_literature_search_tsv()
RETURNS trigger AS $$
BEGIN
NEW.author_names_text := COALESCE(
(SELECT string_agg(value->>'family', ' ')
FROM jsonb_array_elements(NEW.authors)),
''
);
NEW.affiliations_text := COALESCE(
(SELECT string_agg(value->>'affiliation', ' ')
FROM jsonb_array_elements(NEW.authors)),
''
);
NEW.search_tsv := setweight(to_tsvector('english', COALESCE(NEW.title, '')), 'A') ||
setweight(to_tsvector('english', COALESCE(NEW.abstract, '')), 'B') ||
setweight(to_tsvector('simple', COALESCE(NEW.author_names_text, '')), 'A') ||
setweight(to_tsvector('english',
COALESCE(
(SELECT string_agg(value->>'name', ' ')
FROM jsonb_array_elements(NEW.chemical_list)),
'')
), 'C') ||
setweight(to_tsvector('english',
COALESCE(
(SELECT string_agg(value #>> '{}', ' ')
FROM jsonb_array_elements(NEW.gene_symbols)),
'')
), 'C') ||
setweight(to_tsvector('english',
COALESCE(
(SELECT string_agg(value->>'descriptor', ' ')
FROM jsonb_array_elements(NEW.mesh_headings)),
'')
), 'C') ||
setweight(to_tsvector('english',
COALESCE(
(SELECT string_agg(value #>> '{}', ' ')
FROM jsonb_array_elements(NEW.keywords)),
'')
), 'C');
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
"""
_FUNC_OLD = """
CREATE OR REPLACE FUNCTION update_literature_search_tsv()
RETURNS trigger AS $$
BEGIN
NEW.author_names_text := COALESCE(
(SELECT string_agg(value->>'family', ' ')
FROM jsonb_array_elements(NEW.authors)),
''
);
NEW.search_tsv := setweight(to_tsvector('english', COALESCE(NEW.title, '')), 'A') ||
setweight(to_tsvector('english', COALESCE(NEW.abstract, '')), 'B') ||
setweight(to_tsvector('simple', COALESCE(NEW.author_names_text, '')), 'A') ||
setweight(to_tsvector('english',
COALESCE(
(SELECT string_agg(value->>'name', ' ')
FROM jsonb_array_elements(NEW.chemical_list)),
'')
), 'C') ||
setweight(to_tsvector('english',
COALESCE(
(SELECT string_agg(value #>> '{}', ' ')
FROM jsonb_array_elements(NEW.gene_symbols)),
'')
), 'C') ||
setweight(to_tsvector('english',
COALESCE(
(SELECT string_agg(value->>'descriptor', ' ')
FROM jsonb_array_elements(NEW.mesh_headings)),
'')
), 'C') ||
setweight(to_tsvector('english',
COALESCE(
(SELECT string_agg(value #>> '{}', ' ')
FROM jsonb_array_elements(NEW.keywords)),
'')
), 'C');
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
"""
def upgrade() -> None:
# 1. 新增 affiliations_text 列
op.add_column('global_literature', sa.Column('affiliations_text', sa.Text(), nullable=True))
# 2. 创建 trgm 索引(列此刻为空,索引构建秒级;回填时由索引增量维护)
op.execute("""
CREATE INDEX IF NOT EXISTS ix_gl_affiliations_trgm
ON global_literature
USING gin (affiliations_text gin_trgm_ops)
""")
# 3. 更新触发器函数:追加 affiliations_text 赋值(authors 已在触发器 UPDATE OF 列表)
op.execute(_FUNC_NEW)
def downgrade() -> None:
# 3. 还原触发器函数
op.execute(_FUNC_OLD)
# 2. 删除索引
op.execute("DROP INDEX IF EXISTS ix_gl_affiliations_trgm")
# 1. 删除列
op.drop_column('global_literature', 'affiliations_text')
+9
View File
@@ -56,6 +56,7 @@ class GlobalLiterature(Base):
rct_detection: Mapped[dict | None] = mapped_column(JSONB) # RCT 检测结果 {is_rct, confidence, source, evidence}
search_tsv: Mapped[str | None] = mapped_column(TSVECTOR) # 全文检索向量(PG tsvector
author_names_text: Mapped[str | None] = mapped_column(Text) # 从 authors JSONB 提取的 family 文本,由触发器维护,供 trgm 索引
affiliations_text: Mapped[str | None] = mapped_column(Text) # 从 authors JSONB 提取的 affiliation 文本(DISTINCT),由触发器维护,供 trgm 索引
negative_result_details: Mapped[dict | None] = mapped_column(JSONB) # 阴性结果详情
ai_summary: Mapped[dict | None] = mapped_column(JSONB) # AI 摘要 {one_liner, structured, implication}
is_preprint: Mapped[bool] = mapped_column(Boolean, default=False) # 是否为预印本
@@ -95,6 +96,12 @@ class GlobalLiterature(Base):
"is_negative_result", "is_preprint", "journal", "pub_year",
"article_date", "doi", "pmc_id", "language", "citation_status"}),
# 搜索主查询排序索引:ORDER BY pub_date DESC NULLS LAST, id DESCkeyset 分页 + LIMIT)。
# ASC 的 covering index 无法服务 DESC NULLS LAST 顺序;此索引让 LIMITed 主查询走
# Index Scan + Filter,替代 Parallel Seq Scan + Sort89.5 万行 ~40s → <0.2s)。
Index("ix_gl_pub_date_sort",
text("pub_date DESC NULLS LAST, id DESC")),
# 高频筛选 partial index
Index("ix_gl_retracted_true", "retracted", postgresql_where=text("retracted = TRUE")),
Index("ix_gl_is_oa_true", "is_oa", postgresql_where=text("is_oa = TRUE")),
@@ -127,6 +134,8 @@ class GlobalLiterature(Base):
Index("ix_global_literature_journal_trgm", "journal", postgresql_using="gin", postgresql_ops={"journal": "gin_trgm_ops"}),
# Author 搜索:author_names_text 列 + trgm 索引
Index("ix_gl_author_names_trgm", "author_names_text", postgresql_using="gin", postgresql_ops={"author_names_text": "gin_trgm_ops"}),
# Affiliation 搜索:affiliations_text 列 + trgm 索引(替代 jsonb lateral EXISTS
Index("ix_gl_affiliations_trgm", "affiliations_text", postgresql_using="gin", postgresql_ops={"affiliations_text": "gin_trgm_ops"}),
# 常用筛选列 B-tree
Index("ix_gl_doi", "doi"),
Index("ix_gl_journal_iso", "journal_iso"),
+4 -1
View File
@@ -53,10 +53,13 @@ async def expand_atm(db: AsyncSession, query: str):
if not expanded_ids:
return None
# 不用 .distinct()IN (subquery) 语义本身就按值去重,显式 DISTINCT 只会让 PG
# 先物化+排序海量 literature_id'cancer' 这类词展开出 150+ 个 tag,命中的
# literature_id 可达数百万),再对 id IN (百万集合) 扫描。去掉后可直接走
# ix_glt_tag 索引半连接,year_counts 全量扫描和主查询都受益。
return GlobalLiterature.id.in_(
select(GlobalLiteratureTag.literature_id)
.where(GlobalLiteratureTag.tag_id.in_(expanded_ids))
.distinct()
)
+198 -70
View File
@@ -1,5 +1,6 @@
"""高级搜索服务:布尔运算 + 字段限定 + PubMed 查询语法"""
import asyncio
import logging
import re
from datetime import datetime
@@ -27,6 +28,129 @@ def _escape_ilike(s: str) -> str:
return s.replace('\\%', '%').replace('\\_', '_').replace('\\', '\\\\').replace('%', '\\%').replace('_', '\\_')
# 主查询超时(语句级,SET LOCAL 作用于整个事务,之后每条语句都继承该值)
_MAIN_TIMEOUT = "60s"
# year_counts facet 独立短超时:facet 只是辅助直方图,宁可空也不拖垮主查询。
# 前端 axios timeout=15sfacet 预算必须远低于它,否则慢 facet 会让整个请求被客户端取消。
_FACET_TIMEOUT = "8s"
# facetyear_counts 直方图)只对结果集可观的查询有意义:命中过百万时直方图
# 计算本身就是 8s+ 的聚合扫描,还白白拖长请求。总命中数超过该阈值就跳过 facet。
_FACET_TOTAL_MAX = 100_000
# COUNT 只用于 "找到 X 条结果" 展示,keyset 分页不依赖它——给独立短预算,宽查询
# 超时降级 total=0,绝不让 COUNT 拖垮请求(count(*) 对 200 万行命中要 30-40s)。
_COUNT_TIMEOUT = "5s"
# 主查询兜底预算:常规查询(含索引)<1s;此值兜底稀疏匹配的病理查询
# (如 MeSH MAJR + 年份范围,需在百万级日期范围内逐行过滤),超时降级空结果而非 500。
# 预算须与 COUNT 相加 <15s 客户端预算:count 5s + main 6s = 11s。
_MAIN_GUARD_TIMEOUT = "6s"
# P0-F3: 结构化查询(含任何 [MH]/[MAJR]/[TI] 等字段标签)的普通词不做 ATM MeSH 展开。
# 与 PubMed 行为一致:对 "lung neoplasms[MAJR]" PubMed 也不会把 "lung" 自动扩成
# MeSH 树。且 ATM 的 `id IN (tag子查询) OR 全文` 分支会诱导规划器从日期索引逐行
# 过滤——稀疏命中时扫完整日期范围,主查询 6s 超时降级空结果(实测
# "lung neoplasms[MAJR] + 年份" count 5s + main 6s 双超时)。
_STRUCTURED_ATTRS = [
"title_terms", "abstract_terms", "tiab_terms", "author_terms", "journal_terms",
"affiliation_terms", "language_terms", "volume_terms", "issue_terms", "pages_terms",
"lid_terms", "mesh_terms", "majr_terms", "pub_types", "doi_terms", "pmid_terms",
"grant_terms", "subheading_terms", "registry_terms", "substance_terms",
"databank_terms", "pharmaco_terms", "ed_terms", "investigator_terms",
"personal_name_terms", "pubnote_terms", "auid_terms", "cois_terms", "tt_terms",
"sb_terms", "stat_terms", "uid_terms", "ot_terms", "gene_terms", "pmc_terms",
]
_STRUCTURED_DATE_ATTRS = [
"date_from", "date_to", "year_from", "year_to", "edat_from", "edat_to",
"crdt_from", "crdt_to", "mhda_from", "mhda_to", "lr_from", "lr_to",
"dcom_from", "dcom_to", "dep_from", "dep_to",
]
def _has_structured_terms(pp) -> bool:
"""查询是否含任何字段标签/结构化成分(区别于纯自由文本 plain_terms)。"""
if any(getattr(pp, a) for a in _STRUCTURED_ATTRS):
return True
return any(getattr(pp, a) is not None for a in _STRUCTURED_DATE_ATTRS)
async def _query_year_counts(db: AsyncSession, yr_conds: list) -> list[dict]:
"""按年份统计(Results by year facet)。
用独立短超时执行:facet 秒回则用,超时/出错则返回 [],绝不拖垮主查询。
asyncpg 中语句超时会令事务进入 aborted 状态,后续任何语句都报
"current transaction is aborted"——所以失败后必须 rollback 才能继续,
再重设主查询超时(rollback 会把 SET LOCAL 一并清掉)。
返回 [{"year":.., "count":..}],失败返回 []。
"""
await db.execute(text(f"SET LOCAL statement_timeout = '{_FACET_TIMEOUT}'"))
failed = False
year_counts: list[dict] = []
try:
yr_subq = select(GlobalLiterature.pub_year).where(
and_(*yr_conds) if yr_conds else text("TRUE")
).subquery()
year_count_q = select(
yr_subq.c.pub_year, func.count().label("cnt")
).group_by(yr_subq.c.pub_year).order_by(yr_subq.c.pub_year.desc())
year_rows = await db.execute(year_count_q)
year_counts = [
{"year": y, "count": c} for y, c in year_rows if y is not None
]
except asyncio.CancelledError:
raise # 客户端断开,交给上层取消,不吞
except Exception:
failed = True
logger.exception("Year counts query failed")
if failed:
try:
await db.rollback()
except Exception:
logger.exception("Rollback after year counts failure failed")
await db.execute(text(f"SET LOCAL statement_timeout = '{_MAIN_TIMEOUT}'"))
return year_counts
_COUNT_CAP = 100_001 # 宽查询 count 上限:命中超 10 万时停在阈值,避免 count(*) 扫 2M+ 行
async def _safe_count(db: AsyncSession, conds: list) -> int:
"""COUNT(第 1 页 total 展示用)。
宽查询('cancer' 命中 2M+ 行)count(*) 要 30-40s,远超前端 15s 预算。
内层加 LIMIT 上限:扫描到 _COUNT_CAP 行即停止(计划器在 Limit 节点截断),
把宽查询从「5s 超时归零」改为「快速返回 10 万+」。精确小计数不受影响
(命中 < 上限时照常返回真实值)。
仍保留短超时兜底,但超时**绝不返回 0**:宽查询在冷缓存(表 9.6GB,生产
内存有限)下 tag 嵌套循环扫 12 万+ 行 >5s 是常态,返回 0 会让前端显示
"找到 0 条结果" 并隐藏分页(误导)。返回 _COUNT_CAP 表示"结果很多"
facet 因 total>阈值跳过,缓存变暖后下一次即精确。5s 预算保证与 facet(8s)
+ 主查询叠加 <15s 客户端预算。
"""
await db.execute(text(f"SET LOCAL statement_timeout = '{_COUNT_TIMEOUT}'"))
failed = False
total = 0
try:
count_q = select(func.count()).select_from(
select(literal_column("1"))
.select_from(GlobalLiterature)
.where(and_(*conds) if conds else True)
.limit(_COUNT_CAP)
.subquery()
)
total = (await db.execute(count_q)).scalar() or 0
except asyncio.CancelledError:
raise # 客户端断开,交给上层取消,不吞
except Exception:
failed = True
logger.exception("Search count query failed (degrading to cap)")
if failed:
try:
await db.rollback()
except Exception:
logger.exception("Rollback after count failure failed")
total = _COUNT_CAP
await db.execute(text(f"SET LOCAL statement_timeout = '{_MAIN_TIMEOUT}'"))
return total
class AdvancedSearchEngine:
"""PG 高级搜索(ES 就绪后切换 search_service.py"""
@@ -246,8 +370,8 @@ class AdvancedSearchEngine:
"cursor_id": cached.get("cursor_id"),
}
# 30 秒查询超时(放在缓存检查之后,缓存命中不执行)
await db.execute(text("SET LOCAL statement_timeout = '60s'"))
# 查询超时(放在缓存检查之后,缓存命中不执行)
await db.execute(text(f"SET LOCAL statement_timeout = '{_MAIN_TIMEOUT}'"))
# ─── PubMed 语法检测与解析 ───
_pubmed_parsed = None
@@ -464,11 +588,17 @@ class AdvancedSearchEngine:
del conditions[_term_start:]
conditions.append(or_(*_term_conds))
# 年份范围
# 年份范围 —— 用 pub_date 范围(而非 pub_year 整数列),让 ORDER BY pub_date DESC
# 的主查询能走 ix_gl_pub_date_sort 索引的 Index Cond 边界:直接落位到 year_to 内
# 最新行,不必从最新(2024/2025)逐行倒走穿过不匹配年份。
# 数据一致性已核验:pub_year 在范围内 ⟺ (pub_date 在对应范围 OR pub_date IS NULL)
# 反例为 0;改用 pub_date 仅排除 pub_date 为 NULL 的退化行(排序 NULLS LAST
# 深页才可见,对首页结果无影响)。
from datetime import date as _dt_date
if year_from is not None:
conditions.append(GlobalLiterature.pub_year >= year_from)
conditions.append(GlobalLiterature.pub_date >= _dt_date(year_from, 1, 1))
if year_to is not None:
conditions.append(GlobalLiterature.pub_year <= year_to)
conditions.append(GlobalLiterature.pub_date < _dt_date(year_to + 1, 1, 1))
# 具体日期范围(按天搜索)
from datetime import date as dt_date
@@ -623,6 +753,19 @@ class AdvancedSearchEngine:
or has_associated_data or species or sex or age
or medline_only or exclude_preprints)
# ── COUNT 先行(短预算 + 失败回滚)──
# total 驱动 facet 门控,且必须放主查询前执行:若在主查询后跑 COUNT 并失败
# rollback,会把已加载的 ORM items 全部过期,后续逐条懒加载回表(N+1)。
total = 0
if is_first_page:
total = await _safe_count(db, conditions)
# R31: if the main query returns 0 results, year_counts must also be empty
if total == 0 and year_counts:
year_counts = []
# ── facetResults by year)──
# 缓存命中直接用;无筛选走全库缓存;有筛选但命中过多(或 COUNT 失败 total=0
# 跳过直方图——宽查询的聚合要 8s+,宁可不要也不拖垮请求。
if year_counts:
pass # facet 缓存命中
elif not _has_any_filter and not conditions:
@@ -630,40 +773,12 @@ class AdvancedSearchEngine:
if _cached is not None:
year_counts = _cached
else:
try:
yr_conds = conditions[:_yr_before]
yr_subq = select(GlobalLiterature.pub_year).where(
and_(*yr_conds) if yr_conds else text("TRUE")
).subquery()
year_count_q = select(
yr_subq.c.pub_year, func.count().label("cnt")
).group_by(yr_subq.c.pub_year).order_by(yr_subq.c.pub_year.desc())
year_rows = await db.execute(year_count_q)
year_counts = [
{"year": y, "count": c} for y, c in year_rows if y is not None
]
except Exception:
logger.exception("Year counts query failed")
year_counts = []
year_counts = await _query_year_counts(db, conditions[:_yr_before])
await _cache.set("search:year_counts:all", year_counts, ttl=300)
elif conditions:
try:
yr_conds = conditions[:_yr_before]
yr_subq = select(GlobalLiterature.pub_year).where(
and_(*yr_conds)
).subquery()
year_count_q = select(
yr_subq.c.pub_year, func.count().label("cnt")
).group_by(yr_subq.c.pub_year).order_by(yr_subq.c.pub_year.desc())
year_rows = await db.execute(year_count_q)
year_counts = [
{"year": y, "count": c} for y, c in year_rows if y is not None
]
except Exception:
logger.exception("Year counts query failed")
year_counts = []
elif conditions and 0 < total <= _FACET_TOTAL_MAX:
year_counts = await _query_year_counts(db, conditions[:_yr_before])
# 第 1 页结束时统一写 facet 缓存(含 total),此处不重复写入
# 其余情况:total 超阈值 / COUNT 失败 / 无结果 → year_counts 保持 [],直方图略过
# 排序(PubMed 查询时跳过 ts_rank,避免语法标签噪音)
_relevance_query = query
@@ -693,25 +808,42 @@ class AdvancedSearchEngine:
q = q.where(text("FALSE"))
q = q.order_by(*AdvancedSearchEngine._apply_order_by(sort, _relevance_query))
# ── LIMIT page_size+1 探测下一页 + COUNT 第 1 页缓存 ──
# ── LIMIT page_size+1 探测下一页(主查询:独立短预算 + 失败降级空结果)──
# 常规查询(含日期排序索引)<1s;此预算兜底稀疏匹配的病理查询
# MeSH MAJR + 年份范围需在百万级日期范围内逐行过滤),超时返回空结果
# 而非 500。失败后事务已 aborted,必须 rollback 才能继续(后面还有
# load_tags/_get_journal_map 等语句)。
await db.execute(text(f"SET LOCAL statement_timeout = '{_MAIN_GUARD_TIMEOUT}'"))
_main_failed = False
try:
result = await db.execute(q.limit(page_size + 1))
items = result.scalars().all()
except asyncio.CancelledError:
raise # 客户端断开,交给上层取消,不吞
except Exception:
_main_failed = True
logger.exception("Main search query failed (degrading to empty result)")
try:
await db.rollback()
except Exception:
logger.exception("Rollback after main query failure failed")
await db.execute(text(f"SET LOCAL statement_timeout = '{_MAIN_TIMEOUT}'"))
if _main_failed:
return {
"items": [],
"total": 0,
"page": page,
"page_size": page_size,
"has_more": False,
"year_counts": [],
"cursor_val": None,
"cursor_id": None,
}
has_more = len(items) > page_size
items = items[:page_size]
# COUNT 只在第 1 页计算,缓存到 facet key 供后续页复用
total = 0
# COUNT 已在主查询前算好(total);这里统一写 facet 缓存,后续页从缓存读 total
if is_first_page:
count_q = select(func.count()).select_from(
select(literal_column("1"))
.select_from(GlobalLiterature)
.where(and_(*conditions) if conditions else True)
.subquery()
)
total = (await db.execute(count_q)).scalar() or 0
# R31: if the main query returns 0 results, year_counts must also be empty
if total == 0 and year_counts:
year_counts = []
await _cache.set(_facet_cache_key, {"year_counts": year_counts, "total": total}, ttl=300)
else:
# 后续页从 facet 缓存读 total
@@ -913,7 +1045,9 @@ class AdvancedSearchEngine:
neg_conds = [not_(AdvancedSearchEngine._field_condition("all", t.text, t.exact))
for t in pp.plain_terms if t.is_not]
combined_pos_text = " ".join(t.text for t in pp.plain_terms if not t.is_not).strip()
if combined_pos_text and pos_conds and not re.search(r'[一-鿿㐀-䶿豈-﫿]', combined_pos_text):
if (combined_pos_text and pos_conds
and not re.search(r'[一-鿿㐀-䶿豈-﫿]', combined_pos_text)
and not _has_structured_terms(pp)):
from app.services.query_expansion import expand_atm as _expand_atm_inline
try:
atm_cond = await _expand_atm_inline(db, combined_pos_text)
@@ -1870,12 +2004,9 @@ class AdvancedSearchEngine:
GlobalLiterature.journal_iso.ilike(pat),
)
elif field == "affiliation":
# P0-F2: 用 jsonb_array_elements 提取 affiliation 值,避免 JSON 键名假阳性
# P0-F2: affiliations_text 反规范化列(DISTINCT 拼接,trgm 索引),替代 jsonb lateral EXISTS
_pat = _pt()
return text(
"EXISTS (SELECT 1 FROM jsonb_array_elements(global_literature.authors) AS _e "
"WHERE _e->>'affiliation' ILIKE :aff_pat)"
).bindparams(aff_pat=_pat)
return GlobalLiterature.affiliations_text.ilike(_pat)
elif field == "language":
return GlobalLiterature.language.ilike(_pt())
elif field == "volume":
@@ -1902,8 +2033,7 @@ class AdvancedSearchEngine:
GlobalLiterature.abstract.ilike(pat),
GlobalLiterature.journal.ilike(pat),
GlobalLiterature.journal_iso.ilike(pat),
text("EXISTS (SELECT 1 FROM jsonb_array_elements(global_literature.authors) AS _e "
"WHERE _e->>'affiliation' ILIKE :aff_pat)").bindparams(aff_pat=pat),
GlobalLiterature.affiliations_text.ilike(pat),
)
if _wildcard:
# wildcard → ILIKE 右截断(tsvector 不支持 *),多字段覆盖
@@ -1915,8 +2045,7 @@ class AdvancedSearchEngine:
GlobalLiterature.journal_iso.ilike(pat),
cast(GlobalLiterature.pmid, String).ilike(pat),
GlobalLiterature.doi.ilike(pat),
text("EXISTS (SELECT 1 FROM jsonb_array_elements(global_literature.authors) AS _e "
"WHERE _e->>'affiliation' ILIKE :aff_pat)").bindparams(aff_pat=pat),
GlobalLiterature.affiliations_text.ilike(pat),
)
like_val = f"%{_escaped}%"
if "/" in term:
@@ -1924,8 +2053,7 @@ class AdvancedSearchEngine:
return or_(
GlobalLiterature.doi.ilike(_escape_ilike(term)),
GlobalLiterature.doi.ilike(like_val),
text("EXISTS (SELECT 1 FROM jsonb_array_elements(global_literature.authors) AS _e "
"WHERE _e->>'affiliation' ILIKE :aff_pat)").bindparams(aff_pat=pat),
GlobalLiterature.affiliations_text.ilike(pat),
)
return or_(
GlobalLiterature.title.ilike(like_val),
@@ -1935,8 +2063,7 @@ class AdvancedSearchEngine:
GlobalLiterature.author_names_text.ilike(like_val),
GlobalLiterature.journal.ilike(like_val),
GlobalLiterature.journal_iso.ilike(like_val),
text("EXISTS (SELECT 1 FROM jsonb_array_elements(global_literature.authors) AS _e "
"WHERE _e->>'affiliation' ILIKE :aff_pat)").bindparams(aff_pat=pat),
GlobalLiterature.affiliations_text.ilike(pat),
)
# P7-D2: Chinese → ILIKE fallback (tsvector is English-only)
if re.search(r'[一-鿿㐀-䶿豈-﫿]', term):
@@ -1948,8 +2075,7 @@ class AdvancedSearchEngine:
GlobalLiterature.journal_iso.ilike(like_val),
cast(GlobalLiterature.pmid, String).ilike(like_val),
GlobalLiterature.doi.ilike(like_val),
text("EXISTS (SELECT 1 FROM jsonb_array_elements(global_literature.authors) AS _e "
"WHERE _e->>'affiliation' ILIKE :aff_pat)").bindparams(aff_pat=pat),
GlobalLiterature.affiliations_text.ilike(pat),
)
# tsvector 索引主覆盖 title/abstract/author_names/chemicals/genes/mesh/keywords
# journal/journal_iso/affiliation 不在 tsvector 中,以 ILIKE 兜底
@@ -1959,8 +2085,7 @@ class AdvancedSearchEngine:
GlobalLiterature.journal_iso.ilike(like_val),
cast(GlobalLiterature.pmid, String).ilike(like_val),
GlobalLiterature.doi.ilike(like_val),
text("EXISTS (SELECT 1 FROM jsonb_array_elements(global_literature.authors) AS _e "
"WHERE _e->>'affiliation' ILIKE :aff_pat)").bindparams(aff_pat=pat),
GlobalLiterature.affiliations_text.ilike(pat),
)
@staticmethod
@@ -2065,13 +2190,16 @@ class AdvancedSearchEngine:
return None
uids = list(mesh_tag_ids)
# 不用 func.distinctIN (subquery) 语义本身就按值去重,显式 DISTINCT 只让
# PG 多一次物化排序;去掉后可直接走 ix_glt_tag 索引半连接(year_counts 等
# 全量扫描场景更敏感)。
if major_only:
subq = select(func.distinct(GlobalLiteratureTag.literature_id)).where(
subq = select(GlobalLiteratureTag.literature_id).where(
GlobalLiteratureTag.tag_id.in_(uids),
GlobalLiteratureTag.is_major == True,
)
else:
subq = select(func.distinct(GlobalLiteratureTag.literature_id)).where(
subq = select(GlobalLiteratureTag.literature_id).where(
GlobalLiteratureTag.tag_id.in_(uids),
)
return GlobalLiterature.id.in_(subq)
+105
View File
@@ -0,0 +1,105 @@
"""回填 affiliations_text(从 authors JSONB 提取 affiliation 值拼接)
迁移 3355d1f9dbba 只建列/索引/触发器,不做全表回填(单条 UPDATE 持
AccessExclusiveLock 整表,阻塞搜索读)。本脚本分批后台执行:
- 每批 2 万行,批间提交释放锁
- ctid 游标(物理序)顺序推进,每行只读一次,避免 O(n²) 重扫
- 幂等:WHERE affiliations_text IS NULL,可随时重跑/中断续跑
- 串行强制(max_parallel_workers_per_gather=0):并行 SeqScan 不保 ctid 序,
max(ctid) 游标会漏行
- 新插入行由触发器直接维护,脚本只处理存量
- 每 _VACUUM_EVERY 批执行 VACUUM:非 HOT 更新会积累死元组(旧元组+36 索引旧条目
VACUUM 前不释放),全量不清理瞬时占用≈整表大小(27GB)会爆盘。批间 VACUUM
允许并发读写(不阻塞搜索),把瞬时峰值压到 ~1-2GB
用法:
cd backend && python scripts/backfill_affiliations.py # 全量
cd backend && python scripts/backfill_affiliations.py 500000 # 最多处理 N 行后退出
"""
import asyncio
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
import asyncpg
from app.config import Settings
_BATCH = 20_000
_VACUUM_EVERY = 10 # 每 10 批执行一次 VACUUM,回收死元组空间,防止瞬时膨胀爆盘
async def main() -> None:
settings = Settings()
conn = await asyncpg.connect(settings.DATABASE_URL_SYNC)
try:
await conn.execute("SET max_parallel_workers_per_gather = 0")
remaining = await conn.fetchval(
"SELECT count(*) FROM global_literature WHERE affiliations_text IS NULL "
"AND authors IS NOT NULL AND authors != '[]'::jsonb"
)
print(f"待回填: {remaining}")
limit: int | None = None
if len(sys.argv) > 1:
limit = int(sys.argv[1])
print(f"本次最多处理 {limit}")
last: str | None = None
processed = 0
while True:
if limit is not None and processed >= limit:
print(f"达到上限 {limit},退出")
break
if last is None:
rows = await conn.fetch(
"""
SELECT id, ctid::text AS ctid FROM global_literature
WHERE affiliations_text IS NULL
AND authors IS NOT NULL AND authors != '[]'::jsonb
LIMIT $1
""",
_BATCH,
)
else:
rows = await conn.fetch(
"""
SELECT id, ctid::text AS ctid FROM global_literature
WHERE ctid > $1::tid
AND affiliations_text IS NULL
AND authors IS NOT NULL AND authors != '[]'::jsonb
LIMIT $2
""",
last,
_BATCH,
)
if not rows:
break
ids = [r["id"] for r in rows]
last = rows[-1]["ctid"]
await conn.execute(
"""
UPDATE global_literature
SET affiliations_text = COALESCE(
(SELECT string_agg(value->>'affiliation', ' ')
FROM jsonb_array_elements(authors)),
''
)
WHERE id = ANY($1::uuid[])
""",
ids,
)
processed += len(ids)
print(f"已处理 {processed} 行(last ctid {last}", flush=True)
if processed % (_BATCH * _VACUUM_EVERY) == 0:
# VACUUM 并发读写安全,只回收死元组(不阻塞搜索读),压住瞬时膨胀
await conn.execute("VACUUM global_literature")
print(f"已 VACUUM(每 {_VACUUM_EVERY} 批一次)", flush=True)
print(f"回填完成,共处理 {processed}")
finally:
await conn.close()
if __name__ == "__main__":
asyncio.run(main())
+12 -14
View File
@@ -30,7 +30,7 @@ async def test_search_empty_db():
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(db, query="", page=1, page_size=20)
@@ -86,7 +86,7 @@ async def test_search_year_filter():
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(
@@ -101,7 +101,7 @@ async def test_search_boolean_or():
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(15)]
db.execute.return_value = _smart_mock()
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(
@@ -116,7 +116,7 @@ async def test_search_date_range():
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(
@@ -132,10 +132,8 @@ async def test_search_page_size():
"""Custom page_size is respected"""
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
mock_items = MagicMock()
mock_items.scalars.return_value.all.return_value = []
db.execute.side_effect = [_smart_mock() for _ in range(10)] + [mock_items]
db.execute.return_value = _smart_mock()
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(
@@ -154,7 +152,7 @@ async def test_search_retracted_yes():
"""retracted='yes' is now treated as alias for 'only'"""
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(db, retracted="yes", page=1, page_size=20)
assert result["total"] == 0
@@ -165,7 +163,7 @@ async def test_search_retracted_no():
"""retracted='no' filter does not error"""
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(db, retracted="no", page=1, page_size=20)
assert result["total"] == 0
@@ -176,7 +174,7 @@ async def test_search_retracted_only():
"""retracted='only' filter does not error"""
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(db, retracted="only", page=1, page_size=20)
assert result["total"] == 0
@@ -187,7 +185,7 @@ async def test_search_negative_yes():
"""negative_result='yes' is now treated as alias for 'only'"""
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(db, negative_result="yes", page=1, page_size=20)
assert result["total"] == 0
@@ -198,7 +196,7 @@ async def test_search_negative_no():
"""negative_result='no' filter does not error"""
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(db, negative_result="no", page=1, page_size=20)
assert result["total"] == 0
@@ -209,7 +207,7 @@ async def test_search_negative_only():
"""negative_result='only' filter does not error"""
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(db, negative_result="only", page=1, page_size=20)
assert result["total"] == 0
@@ -220,7 +218,7 @@ async def test_search_retracted_and_negative():
"""Combined retracted + negative_result filters do not error"""
from app.services.search_engine import AdvancedSearchEngine
db = AsyncMock()
db.execute.side_effect = [_smart_mock() for _ in range(10)]
db.execute.return_value = _smart_mock() # 无界 mock:搜索路径语句数会随预算/降级逻辑增长,固定 side_effect 列表易被打爆
with patch("app.services.tag_loader.load_tags_for_literature", AsyncMock(return_value={})):
result = await AdvancedSearchEngine.search(
db, retracted="yes", negative_result="only", page=1, page_size=20
+12 -1
View File
@@ -179,6 +179,8 @@ CREATE TABLE global_literature (
is_negative_result BOOLEAN NOT NULL DEFAULT FALSE, -- 是否为阴性结果(AI 判定,非 PubMed 原始字段)
rct_detection JSONB, -- RCT 自动检测结果 {is_rct, confidence, source}AI,非 PubMed 原始字段)
search_tsv TSVECTOR, -- 全文检索向量(内部 PG tsvector,非 PubMed 字段)
author_names_text TEXT, -- 作者姓名字符串(从 authors.family 拼接,由触发器维护,供 trgm 索引 ILIKE
affiliations_text TEXT, -- 机构字符串(从 authors.affiliation 拼接,由触发器维护,供 trgm 索引 ILIKE
negative_result_details JSONB, -- 阴性结果详情(AI 抽取,非 PubMed 原始字段)
ai_summary JSONB, -- AI 摘要(内部,非 PubMed 字段)
chemical_list JSONB NOT NULL DEFAULT '[]', -- 化学物质列表 【PubMed: ChemicalList/Chemical】
@@ -235,12 +237,21 @@ CREATE INDEX ix_gl_is_preprint_true ON global_literature(is_preprint) WHERE is_p
-- ═══ GIN 反范式标签数组索引(标签筛选免 JOIN global_literature_tags ═══
CREATE INDEX ix_gl_tag_ids_gin ON global_literature USING gin(tag_ids);
CREATE INDEX ix_global_literature_journal_iso_trgm ON global_literature USING gin(journal_iso gin_trgm_ops); -- PubMed [TA] ILIKE 兜底
-- ═══ 反范式文本列 trgm 索引(ILIKE 走索引,避免 jsonb lateral EXISTS ═══
CREATE INDEX ix_gl_author_names_trgm ON global_literature USING gin(author_names_text gin_trgm_ops); -- 作者 [AU] ILIKE
CREATE INDEX ix_gl_affiliations_trgm ON global_literature USING gin(affiliations_text gin_trgm_ops); -- 机构 [AD] ILIKE
-- ═══ 排序索引(pub_date DESC NULLS LAST, id DESC,高级搜索按时间排序) ═══
CREATE INDEX ix_gl_pub_date_sort ON global_literature(pub_date DESC NULLS LAST, id DESC);
```
> 预估值:约 600 万条(PubMed 中肿瘤相关的历史累积),单表不需要分区
**全文搜索:**
> `search_tsv``TSVECTOR` 类型,GIN 索引 `ix_gl_search_tsv`)由 PG 触发器 `trg_global_literature_tsv` 自动维护。使用 `setweight()` 区分字段权重:**title=A 权重、abstract=B 权重、authors.family=A 权重**。所有搜索(普通搜索 `field="all"` 和高级搜索字段选择中的作者/机构)均走 tsvector `@@` `plainto_tsquery()`ILIKE 仅作为兜底(NULL tsvector 记录)。触发器在 INSERT 或 UPDATE title/abstract/authors 时触发重建。千万级无压力。
> `search_tsv``TSVECTOR` 类型,GIN 索引 `ix_gl_search_tsv`)由 PG 触发器 `trg_global_literature_tsv` 自动维护。使用 `setweight()` 区分字段权重:**title=A 权重、abstract=B 权重、author_names_text=A 权重**。所有搜索(普通搜索 `field="all"` 和高级搜索字段选择中的作者/机构)均走 tsvector `@@` `plainto_tsquery()`ILIKE 仅作为兜底(NULL tsvector 记录)。触发器在 INSERT 或 UPDATE title/abstract/authors/chemical_list/gene_symbols/mesh_headings/keywords 时触发重建。千万级无压力。
>
> **反规范化文本列:** `author_names_text`authors.family 拼接)和 `affiliations_text`authors.affiliation 拼接)由同一触发器维护,配套 `gin_trgm_ops` 索引让作者 [AU]/机构 [AD] 的 ILIKE 搜索走索引——替代逐行 `jsonb_array_elements(authors)` 的 lateral EXISTS(病理查询下慢 ~40x)。`affiliations_text` 迁移 `3355d1f9dbba` 只建列/索引/触发器,存量回填由 `scripts/backfill_affiliations.py` 分批后台执行(每批 2 万行、批间提交,幂等 `WHERE affiliations_text IS NULL`);新插入/更新的行由触发器直接维护。
>
> > **GIN 索引重建说明:** 迁移 `f80dca5baa02`add_pico_column_to_review_literatures)的 `upgrade()` 误将 `ix_gl_search_tsv` 删除后未重建,导致所有 tsvector 搜索走全表扫描。迁移 `0314f4d28728`rebuild_search_tsv_gin_index)已修复,紧跟在 `e341edea85e2` 之后。