如何用 Elasticsearch 设计一个全站搜索引擎
多数据源索引、CDC 同步、权限感知搜索、热词与联想输入——一份完整的 Elasticsearch 架构指南。
构建一个全站搜索引擎听起来是条走过千百遍的老路——直到你面对真实世界的种种约束:多个数据源(既有关系型数据库又有第三方 API)、格式各异的文档(HTML、PDF、Word、Excel、PowerPoint)、细粒度的权限过滤、带时间衰减的热词排行,以及联想输入建议。本文将走过一套以 Elasticsearch 8.x 为核心搜索引擎、涵盖上述全部关注点的生产级设计。
需求概览
搜索引擎必须支持四个核心功能:
- 关键词搜索——对标题和正文内容做全文检索,支持命中高亮、可配置的来源优先级加权、相关性与时间排序,以及按用户维度的权限过滤。
- 多源混合排序——来自我方 Elasticsearch 索引的结果必须与第三方 API 的结果合并,按统一的相关性得分排序,并在两个来源间实现正确的分页。
- 热词排行——一份每日更新的热门搜索词榜单,采用时间衰减公式,让过时的词自然淡出。
- 联想输入建议——由 Elasticsearch Completion Suggester 驱动的边输入边补全。
架构
系统分为两大层:
搜索服务(Search Service) 负责所有查询期逻辑:客户端调用的搜索 API、权限感知过滤、热词获取、联想输入建议,以及合并 Elasticsearch 与外部 API 结果的混合排序引擎。
数据同步管道(Data Synchronization Pipeline) 通过解析 binlog 的 CDC(变更数据捕获)技术,让 Elasticsearch 与权威数据源 MySQL 保持近实时同步。在写入索引前,schema 转换器会将原始变更事件转换成 Elasticsearch 文档格式。
搜索行为日志(查询词、时间戳、用户 ID)会写入 MySQL,供每日批处理作业用于计算热词和联想候选词。
技术选型
3.1 搜索引擎:Elasticsearch 8.x
Elasticsearch 仍是最成熟、维护最活跃的开源全文搜索引擎。对于云上部署,Amazon OpenSearch Service 提供了一个与 Elasticsearch API 兼容的托管替代方案,免去了集群管理的运维负担。
关键配置要点:
- 启用 IK 分词插件 以支持 CJK(中日韩)分词。如果需要细粒度的中文切分,在索引时使用
ik_max_word,在搜索时使用ik_smart。 - 谨慎规划堆内存。重型分词插件在小实例上可能导致 OOM——生产环境请预留至少 8 GB 的 JVM 堆内存。
3.2 文档抽取
以二进制文件形式存储的文档(PDF、Word、Excel、PowerPoint)需要先转换成纯文本才能被索引。主要选项如下:
| 工具 | 方式 | 权衡 |
|---|---|---|
| Apache Tika | Java 库,格式支持最广 | 需要 JVM;复杂排版可能丢失保真度 |
| Ingest Attachment | 封装了 Tika 的 ES 插件 | 集成度高,但运行在 ES 节点内部 |
| FsCrawler | 独立的文件系统爬虫 | 适合批处理;不适合流式场景 |
| 云 API | AWS Textract、Azure AI Document Intelligence | 按页计费;对扫描件精度最高 |
对于流式架构,推荐的做法是在写入 Elasticsearch 之前,于 schema 转换器内部调用文档抽取服务(Tika 或云 API)。这样能让抽取逻辑与应用及 ES 本身都解耦。
3.3 数据同步:基于 binlog 的 CDC
将 MySQL 数据同步到 Elasticsearch,常见有四种模式:
| 模式 | 优点 | 缺点 |
|---|---|---|
| 同步双写 | 延迟最低 | 强耦合;有部分失败风险 |
| 异步双写(经 MQ) | 写入解耦 | 应用必须自行发布事件 |
| 定期批量 ETL | 应用零改动 | 延迟高;对 MySQL 有轮询压力 |
| Binlog CDC | 近实时;应用零改动;一致性好 | 需要 CDC 工具 |
Binlog CDC 是推荐方案,因为它对应用透明,提供近实时的延迟,并保证每一次已提交的变更都被恰好捕获一次。
CDC 工具选型
使用最广泛的两个开源 CDC 连接器是:
- Canal(阿里巴巴)——一个成熟的 Java 工具,通过伪装成 MySQL 从库来接收 binlog 事件。它在中文技术生态中拥有庞大的用户群,并支持 Kafka、RocketMQ 以及自定义下游适配器。
- Debezium(Red Hat)——一个基于 Kafka Connect 的 CDC 平台,支持 MySQL、PostgreSQL、MongoDB 等众多数据库。它提供更丰富的事件格式(变更前/后快照)、内置的 schema 演进处理,是 Kafka 为中心架构中的标准选择。
对于 AWS 原生部署,AWS Database Migration Service(DMS) 只需极少配置,即可将 CDC 事件从 RDS MySQL 流式传输到 Amazon MSK(Kafka)、Amazon OpenSearch 或 S3。
管道的工作流程如下:
- MySQL 为每一笔已提交的事务写入 binlog 事件。
- CDC 连接器(Canal、Debezium 或 AWS DMS)伪装成 MySQL 从库来读取 binlog 流。
- 变更事件被发布到一个 Kafka topic,每次行变更对应一个事件。
- schema 转换器(Logstash、自定义服务或 Kafka Streams)消费这些事件,将数据库列映射为 Elasticsearch 字段,可选地抽取文档文本,并生成一个包含 ES 就绪 JSON 文档的新 Kafka topic。
- Logstash(或自定义消费者)从输出 topic 读取数据,调用 Elasticsearch Bulk API 将文档写入索引。
索引设计
4.1 全文搜索索引模板
Elasticsearch 8.x 使用 可组合索引模板(_index_template),而非遗留的 _template API。以下是全文搜索模板:
PUT _index_template/template_fulltext
{
"index_patterns": ["fulltext-*"],
"template": {
"settings": {
"number_of_shards": 1,
"number_of_replicas": 1,
"analysis": {
"analyzer": {
"ik_index_analyzer": {
"type": "custom",
"tokenizer": "ik_max_word"
},
"ik_search_analyzer": {
"type": "custom",
"tokenizer": "ik_smart"
}
}
}
},
"mappings": {
"properties": {
"title": {
"type": "text",
"analyzer": "ik_index_analyzer",
"search_analyzer": "ik_search_analyzer"
},
"summary": {
"type": "text",
"analyzer": "ik_index_analyzer",
"search_analyzer": "ik_search_analyzer"
},
"content": {
"type": "text",
"analyzer": "ik_index_analyzer",
"search_analyzer": "ik_search_analyzer"
},
"author": {
"type": "keyword"
},
"document_type": {
"type": "keyword"
},
"url": {
"type": "keyword"
},
"publish_date": {
"type": "date"
},
"update_date": {
"type": "date"
},
"privilege": {
"properties": {
"data": {
"type": "nested",
"properties": {
"type": {
"type": "keyword"
},
"id": {
"type": "keyword"
}
}
}
}
}
}
}
}
}
注意: 遗留的
PUT _template/template_nameAPI 在 ES 7.8+ 中已废弃,并在 ES 9 中被移除。对于 ES 8.x 及以上版本,请始终使用带可组合模板的PUT _index_template/template_name。
4.2 建议索引模板
联想输入索引使用 completion 字段类型:
PUT _index_template/template_suggest
{
"index_patterns": ["suggest-*"],
"template": {
"settings": {
"number_of_shards": 1
},
"mappings": {
"properties": {
"suggest": {
"type": "completion"
},
"weight": {
"type": "integer"
}
}
}
}
}
权限感知搜索
在企业搜索中,不同用户能看到的文档各不相同。有两种策略:
| 策略 | 何时过滤 | 权衡 |
|---|---|---|
| 检索前过滤 | 查询期(ES filter 子句) | 性能最佳;查询更复杂 |
| 检索后过滤 | ES 返回结果之后 | 查询更简单;浪费检索预算 |
追求高性能搜索时,强烈推荐 检索前过滤。我们将权限数据以 nested 字段的形式直接嵌入文档:
{
"title": "Q3 Financial Report",
"content": "...",
"privilege": {
"data": [
{ "type": "staff", "id": "user-1234" },
{ "type": "department", "id": "dept-finance" },
{ "type": "department", "id": "dept-executive" }
]
}
}
privilege.data 中的每一项代表一个被允许查看该文档的实体(用户或部门)。只要用户匹配到任意一个权限项,就应当能看到该文档——要么其自身的用户 ID 出现在某个 staff 项中,或者 其部门 ID 出现在某个 department 项中。
权限过滤查询(修正版)
正确的查询在 staff 与 department 两个 nested 查询之间使用 bool.should(OR):
GET /fulltext-*/_search
{
"query": {
"bool": {
"must": [
{
"multi_match": {
"query": "financial report",
"fields": ["title^3", "summary^2", "content"]
}
}
],
"filter": [
{
"bool": {
"should": [
{
"nested": {
"path": "privilege.data",
"query": {
"bool": {
"must": [
{ "term": { "privilege.data.type": "staff" } },
{ "term": { "privilege.data.id": "user-1234" } }
]
}
}
}
},
{
"nested": {
"path": "privilege.data",
"query": {
"bool": {
"must": [
{ "term": { "privilege.data.type": "department" } },
{ "term": { "privilege.data.id": "dept-finance" } }
]
}
}
}
}
],
"minimum_should_match": 1
}
}
]
}
},
"highlight": {
"fields": {
"title": {},
"content": { "fragment_size": 200 }
}
}
}
踩坑修正说明: 一个常见错误是把两个 nested 查询作为
filter数组中的两个独立项,这会施加 AND 语义——即用户必须同时匹配一个 staff 项和一个 department 项。这会导致大多数文档被过滤掉。正确的做法是用带minimum_should_match: 1的bool.should把它们包起来,这样匹配任一条件即可。
多源混合排序与分页
这是整个系统在架构上最有意思的部分。当搜索结果同时来自 Elasticsearch(用 BM25 打分)和第三方 API(用其自有算法打分)时,我们需要:
- 并发查询两个来源——使用异步/并行调用,避免串行带来的延迟。
- 归一化得分——要么把第三方得分重新缩放到 ES 的得分区间,要么为每个来源施加可配置的权重(例如 ES 结果乘以 1.2 倍)。
- 合并并排序——按得分降序交错排列结果,产出单一、统一的页面。
- 跟踪各来源的偏移量——由于每页从各来源消费的条目数不同,我们需要双游标。
6.1 双偏移量分页算法
当结果来自两个独立来源时,标准的基于偏移量的分页(from + size)会失效。解决方案是跟踪两个独立的偏移量——每个来源一个——并让合并过程来决定每页中各来源各出现多少条目。
# Pseudocode for the mixed-ranking search endpoint
async def search(query: str, es_offset: int, api_offset: int, page_size: int):
# 1. Fetch from both sources in parallel
es_results, api_results = await asyncio.gather(
search_elasticsearch(query, offset=es_offset, limit=page_size),
search_third_party(query, offset=api_offset, limit=page_size),
)
# 2. Merge by score (descending)
merged = []
es_used, api_used = 0, 0
es_idx, api_idx = 0, 0
while len(merged) < page_size:
es_item = es_results[es_idx] if es_idx < len(es_results) else None
api_item = api_results[api_idx] if api_idx < len(api_results) else None
if es_item is None and api_item is None:
break
if api_item is None or (es_item and es_item.score >= api_item.score):
merged.append(es_item)
es_idx += 1
es_used += 1
else:
merged.append(api_item)
api_idx += 1
api_used += 1
return {
"results": merged,
"es_offset": es_offset,
"api_offset": api_offset,
"es_used": es_used,
"api_used": api_used,
"es_has_next": es_results.has_more,
"api_has_next": api_results.has_more,
}
6.2 客户端分页状态
前端维护一个偏移量快照数组——每访问过一页对应一个条目:
interface PageState {
esOffset: number;
apiOffset: number;
}
// Initialize
const pageHistory: PageState[] = [{ esOffset: 0, apiOffset: 0 }];
// After receiving page results:
function onNextPage(response: SearchResponse) {
const nextState: PageState = {
esOffset: response.es_offset + response.es_used,
apiOffset: response.api_offset + response.api_used,
};
pageHistory.push(nextState);
}
// Go to previous page:
function onPrevPage() {
pageHistory.pop(); // remove current
const prev = pageHistory[pageHistory.length - 1];
// re-fetch with prev.esOffset, prev.apiOffset
}
以 page_size = 20 为例的推演:
| 页码 | 请求 | 响应 | 状态 |
|---|---|---|---|
| 1 | esOffset=0, apiOffset=0 | es_used=7, api_used=13 | [{0,0}] |
| 2 | esOffset=7, apiOffset=13 | es_used=12, api_used=8 | [{0,0}, {7,13}] |
| 返回第 1 页 | 弹出 {7,13} → 重新拉取 {0,0} | 同第 1 页 | [{0,0}] |
这种做法牺牲了随机跳页的能力(你无法直接跳到第 5 页),但它妥善处理了合并两个独立结果流这一根本性复杂问题,并做到正确、确定的分页。对大多数搜索 UI 而言,上一页/下一页导航已经足够。
带时间衰减的热词排行
一个朴素的热词实现只是简单地统计固定窗口内(例如最近 30 天)的搜索频次。但这会带来一个问题:一个 10 天前被搜过 101 次的关键词,会排在一个昨天被搜过 100 次的关键词前面,尽管后者显然此刻”更热”。
解决方案是一个线性时间衰减公式,赋予近期搜索更高的权重。
7.1 衰减公式
给定:
- T = 时间窗口大小(例如 30 天)
- c_i = 第 i 天的搜索次数,其中 i = 0 表示昨天,i = T-1 表示最早的一天
- w = 基础权重(默认 1.0)
某关键词的热度得分为:
hot = (w / T) × Σ (T - i) × cᵢ for i = 0 to T-1
展开为 T = 30 时:
hot = (30 * c_0 + 29 * c_1 + 28 * c_2 + ... + 1 * c_29) / 30
工作原理: 昨天的搜索(c_0)乘以 30,前天(c_1)乘以 29,以此类推。30 天前的搜索(c_29)只乘以 1。这形成了一条平滑的线性衰减曲线,近期活跃度的权重最高可达窗口边缘活跃度的 30 倍。
7.2 实现
from datetime import datetime, timedelta
from collections import defaultdict
def compute_hot_keywords(
search_logs: list[dict], # [{"keyword": str, "date": date}, ...]
window_days: int = 30,
top_n: int = 50,
base_weight: float = 1.0,
) -> list[dict]:
"""Compute hot keyword scores with linear time decay."""
today = datetime.utcnow().date()
scores = defaultdict(float)
for log in search_logs:
keyword = log["keyword"]
days_ago = (today - log["date"]).days
if days_ago < 0 or days_ago >= window_days:
continue
# Linear decay: recent days get higher weight
decay_factor = window_days - days_ago
scores[keyword] += decay_factor * base_weight / window_days
# Sort by score descending, take top N
ranked = sorted(scores.items(), key=lambda x: -x[1])[:top_n]
return [{"keyword": k, "score": round(s, 2)} for k, s in ranked]
每日 cron 作业计算这些得分,把前 N 名存入 MySQL(供人工调整——编辑可以置顶、删除或重排词条),并将最终列表缓存到 Redis,供搜索 API 亚毫秒级读取。
7.3 扩展
- 同义词合并: 在打分前先归一化同义词,让 “ES”、“Elasticsearch” 和 “elastic search” 被算作同一个关键词。对大多数场景来说,一张简单的别名映射表或一个编辑距离阈值就够用了。
- 突发检测: 为那些当日次数超过滚动均值 2 倍的关键词加一个乘数——这能更积极地凸显突然的峰值。
- 指数衰减变体: 把
(T - i)替换为e^{-lambda * i},得到更陡峭的下降。线性公式更易于解释和调优,但指数衰减在压制”老而频繁”的词条方面表现更好。
用 Completion Suggester 实现联想输入建议
Elasticsearch 的 Completion Suggester 是一种专为前缀自动补全优化的数据结构(FST——有限状态转换器)。它完全驻留在内存中,能以亚毫秒级返回结果。
8.1 索引建议词
每日批处理作业从最近 90 天中计算出排名前 1000+ 的搜索词,并写入建议索引:
POST suggest-v1/_doc
{
"suggest": {
"input": ["elasticsearch", "elastic search", "ES"],
"weight": 85
}
}
input 数组允许多种书写形式(包括常见拼写错误)映射到同一个建议。weight 字段控制排序——权重越高的建议越靠前。
8.2 查询建议词
GET suggest-v1/_search
{
"suggest": {
"keyword-suggest": {
"prefix": "elast",
"completion": {
"field": "suggest",
"size": 10,
"skip_duplicates": true
}
}
}
}
这会返回 input 以 “elast” 开头的前 10 个建议,去重后按权重排序。
组合起来:搜索 API
以下是把所有部分串联起来的搜索接口的高层视图:
@app.get("/api/search")
async def search(
q: str,
sort_by: str = "relevance", # "relevance" or "date"
es_offset: int = 0,
api_offset: int = 0,
page_size: int = 20,
user: User = Depends(get_current_user),
):
# 1. Build the ES query with permission filter
es_query = build_search_query(
keyword=q,
user_id=user.id,
department_id=user.department_id,
sort_by=sort_by,
)
# 2. Query ES and third-party API concurrently
es_results, api_results = await asyncio.gather(
es_client.search(
index="fulltext-*",
body=es_query,
from_=es_offset,
size=page_size,
),
third_party_client.search(q, offset=api_offset, limit=page_size),
)
# 3. Merge results by score
merged = merge_and_rank(es_results, api_results, page_size)
# 4. Log the search event (async, non-blocking)
asyncio.create_task(log_search_event(user.id, q))
return merged
@app.get("/api/search/hot")
async def hot_keywords():
"""Return cached hot keywords (refreshed daily by cron)."""
cached = await redis.get("search:hot_keywords")
if cached:
return json.loads(cached)
return await refresh_hot_keywords_from_db()
@app.get("/api/search/suggest")
async def suggest(prefix: str):
"""Typeahead suggestions via ES Completion Suggester."""
result = await es_client.search(
index="suggest-*",
body={
"suggest": {
"keyword-suggest": {
"prefix": prefix,
"completion": {
"field": "suggest",
"size": 10,
"skip_duplicates": True,
}
}
}
}
)
suggestions = [
opt["text"]
for opt in result["suggest"]["keyword-suggest"][0]["options"]
]
return {"suggestions": suggestions}
未来优化方向
数据质量
- 文档去重——用 SimHash 或 MinHash 检测近似重复的文档,并在索引时合并它们。
- 内容清洗——在索引前,从 HTML 文档中剥离样板内容(导航、页脚、广告)。
- 来源级权重配置——允许管理员在不重新部署的情况下设置按索引或按来源的加权因子。
搜索质量
- 自定义词典管理——用来自热词表和人工整理的领域专有词条来扩充 IK 分词器的词典。
- 同义词扩展——配置 ES 同义词词元过滤器,让 “K8s” 匹配 “Kubernetes”。
- 拼写纠错——使用 Phrase Suggester 或专门的拼写检查层来处理错别字。
- 点击率(CTR)跟踪——记录用户实际点击了哪些结果,并通过
function_score包装器把这一信号反哺到相关性打分中。
基础设施
- 索引生命周期管理(ILM)——自动滚动、收缩和删除旧索引,以控制存储成本。
- 搜索相关性测试——在调优分词器和打分时,使用 ES Ranking Evaluation API 运行自动化的相关性基准测试。
- 可观测性——把 p50/p95/p99 查询延迟、零结果率和建议采纳率作为关键的搜索健康指标进行跟踪。
参考资料
参考资料
- Elasticsearch reference — Elastic
- IK Analysis plugin for Chinese — GitHub
- Elasticsearch text analyzers — Elastic