如何用 Elasticsearch 设计一个全站搜索引擎

多数据源索引、CDC 同步、权限感知搜索、热词与联想输入——一份完整的 Elasticsearch 架构指南。

zhuermu··20 分钟
ElasticsearchSearch EngineCDCCanalFull-Text SearchSystem Design

构建一个全站搜索引擎听起来是条走过千百遍的老路——直到你面对真实世界的种种约束:多个数据源(既有关系型数据库又有第三方 API)、格式各异的文档(HTML、PDF、Word、Excel、PowerPoint)、细粒度的权限过滤、带时间衰减的热词排行,以及联想输入建议。本文将走过一套以 Elasticsearch 8.x 为核心搜索引擎、涵盖上述全部关注点的生产级设计。


需求概览

搜索引擎必须支持四个核心功能:

  1. 关键词搜索——对标题和正文内容做全文检索,支持命中高亮、可配置的来源优先级加权、相关性与时间排序,以及按用户维度的权限过滤。
  2. 多源混合排序——来自我方 Elasticsearch 索引的结果必须与第三方 API 的结果合并,按统一的相关性得分排序,并在两个来源间实现正确的分页。
  3. 热词排行——一份每日更新的热门搜索词榜单,采用时间衰减公式,让过时的词自然淡出。
  4. 联想输入建议——由 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 TikaJava 库,格式支持最广需要 JVM;复杂排版可能丢失保真度
Ingest Attachment封装了 Tika 的 ES 插件集成度高,但运行在 ES 节点内部
FsCrawler独立的文件系统爬虫适合批处理;不适合流式场景
云 APIAWS 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。

数据同步

管道的工作流程如下:

  1. MySQL 为每一笔已提交的事务写入 binlog 事件。
  2. CDC 连接器(Canal、Debezium 或 AWS DMS)伪装成 MySQL 从库来读取 binlog 流。
  3. 变更事件被发布到一个 Kafka topic,每次行变更对应一个事件。
  4. schema 转换器(Logstash、自定义服务或 Kafka Streams)消费这些事件,将数据库列映射为 Elasticsearch 字段,可选地抽取文档文本,并生成一个包含 ES 就绪 JSON 文档的新 Kafka topic。
  5. 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_name API 在 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: 1bool.should 把它们包起来,这样匹配任一条件即可。


多源混合排序与分页

这是整个系统在架构上最有意思的部分。当搜索结果同时来自 Elasticsearch(用 BM25 打分)和第三方 API(用其自有算法打分)时,我们需要:

  1. 并发查询两个来源——使用异步/并行调用,避免串行带来的延迟。
  2. 归一化得分——要么把第三方得分重新缩放到 ES 的得分区间,要么为每个来源施加可配置的权重(例如 ES 结果乘以 1.2 倍)。
  3. 合并并排序——按得分降序交错排列结果,产出单一、统一的页面。
  4. 跟踪各来源的偏移量——由于每页从各来源消费的条目数不同,我们需要双游标。

混合排序

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 为例的推演

页码请求响应状态
1esOffset=0, apiOffset=0es_used=7, api_used=13[{0,0}]
2esOffset=7, apiOffset=13es_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 查询延迟、零结果率和建议采纳率作为关键的搜索健康指标进行跟踪。

参考资料

参考资料

  1. Elasticsearch reference — Elastic
  2. IK Analysis plugin for Chinese — GitHub
  3. Elasticsearch text analyzers — Elastic