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

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

- 作者: zhuermu
- 发布: 2024-03-10
- 网页版: https://zhuermu.com/blog/full-site-search-engine-elasticsearch/
- 首发于: https://blog.csdn.net/qq258513813/article/details/136712280

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

---

## 1. 需求概览

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

1. **关键词搜索**——对标题和正文内容做全文检索，支持命中高亮、可配置的来源优先级加权、相关性与时间排序，以及按用户维度的权限过滤。
2. **多源混合排序**——来自我方 Elasticsearch 索引的结果必须与第三方 API 的结果合并，按统一的相关性得分排序，并在两个来源间实现正确的分页。
3. **热词排行**——一份每日更新的热门搜索词榜单，采用时间衰减公式，让过时的词自然淡出。
4. **联想输入建议**——由 Elasticsearch Completion Suggester 驱动的边输入边补全。

---

## 2. 架构

![架构图](/images/blog/full-site-search-engine-elasticsearch/architecture.svg)

系统分为两大层：

**搜索服务（Search Service）** 负责所有查询期逻辑：客户端调用的搜索 API、权限感知过滤、热词获取、联想输入建议，以及合并 Elasticsearch 与外部 API 结果的混合排序引擎。

**数据同步管道（Data Synchronization Pipeline）** 通过解析 binlog 的 CDC（变更数据捕获）技术，让 Elasticsearch 与权威数据源 MySQL 保持近实时同步。在写入索引前，schema 转换器会将原始变更事件转换成 Elasticsearch 文档格式。

搜索行为日志（查询词、时间戳、用户 ID）会写入 MySQL，供每日批处理作业用于计算热词和联想候选词。

---

## 3. 技术选型

### 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。

![数据同步](/images/blog/full-site-search-engine-elasticsearch/data-sync.svg)

管道的工作流程如下：

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. 索引设计

### 4.1 全文搜索索引模板

Elasticsearch 8.x 使用 **可组合索引模板**（`_index_template`），而非遗留的 `_template` API。以下是全文搜索模板：

```json
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` 字段类型：

```json
PUT _index_template/template_suggest
{
  "index_patterns": ["suggest-*"],
  "template": {
    "settings": {
      "number_of_shards": 1
    },
    "mappings": {
      "properties": {
        "suggest": {
          "type": "completion"
        },
        "weight": {
          "type": "integer"
        }
      }
    }
  }
}
```

---

## 5. 权限感知搜索

在企业搜索中，不同用户能看到的文档各不相同。有两种策略：

| 策略 | 何时过滤 | 权衡 |
|----------|---------------|-----------|
| **检索前过滤** | 查询期（ES `filter` 子句） | 性能最佳；查询更复杂 |
| 检索后过滤 | ES 返回结果之后 | 查询更简单；浪费检索预算 |

追求高性能搜索时，强烈推荐 **检索前过滤**。我们将权限数据以 `nested` 字段的形式直接嵌入文档：

```json
{
  "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）：

```json
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` 把它们包起来，这样匹配*任一*条件即可。

---

## 6. 多源混合排序与分页

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

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

![混合排序](/images/blog/full-site-search-engine-elasticsearch/mixed-ranking.svg)

### 6.1 双偏移量分页算法

当结果来自两个独立来源时，标准的基于偏移量的分页（`from` + `size`）会失效。解决方案是跟踪**两个独立的偏移量**——每个来源一个——并让合并过程来决定每页中各来源各出现多少条目。

```python
# 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 客户端分页状态

前端维护一个偏移量快照数组——每访问过一页对应一个条目：

```typescript
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 而言，上一页/下一页导航已经足够。

---

## 7. 带时间衰减的热词排行

一个朴素的热词实现只是简单地统计固定窗口内（例如最近 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 实现

```python
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}`，得到更陡峭的下降。线性公式更易于解释和调优，但指数衰减在压制"老而频繁"的词条方面表现更好。

---

## 8. 用 Completion Suggester 实现联想输入建议

Elasticsearch 的 Completion Suggester 是一种专为前缀自动补全优化的数据结构（FST——有限状态转换器）。它完全驻留在内存中，能以亚毫秒级返回结果。

### 8.1 索引建议词

每日批处理作业从最近 90 天中计算出排名前 1000+ 的搜索词，并写入建议索引：

```json
POST suggest-v1/_doc
{
  "suggest": {
    "input": ["elasticsearch", "elastic search", "ES"],
    "weight": 85
  }
}
```

`input` 数组允许多种书写形式（包括常见拼写错误）映射到同一个建议。`weight` 字段控制排序——权重越高的建议越靠前。

### 8.2 查询建议词

```json
GET suggest-v1/_search
{
  "suggest": {
    "keyword-suggest": {
      "prefix": "elast",
      "completion": {
        "field": "suggest",
        "size": 10,
        "skip_duplicates": true
      }
    }
  }
}
```

这会返回 input 以 "elast" 开头的前 10 个建议，去重后按权重排序。

---

## 9. 组合起来：搜索 API

以下是把所有部分串联起来的搜索接口的高层视图：

```python
@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}
```

---

## 10. 未来优化方向

### 数据质量
- **文档去重**——用 SimHash 或 MinHash 检测近似重复的文档，并在索引时合并它们。
- **内容清洗**——在索引前，从 HTML 文档中剥离样板内容（导航、页脚、广告）。
- **来源级权重配置**——允许管理员在不重新部署的情况下设置按索引或按来源的加权因子。

### 搜索质量
- **自定义词典管理**——用来自热词表和人工整理的领域专有词条来扩充 IK 分词器的词典。
- **同义词扩展**——配置 ES 同义词词元过滤器，让 "K8s" 匹配 "Kubernetes"。
- **拼写纠错**——使用 Phrase Suggester 或专门的拼写检查层来处理错别字。
- **点击率（CTR）跟踪**——记录用户实际点击了哪些结果，并通过 `function_score` 包装器把这一信号反哺到相关性打分中。

### 基础设施
- **索引生命周期管理（ILM）**——自动滚动、收缩和删除旧索引，以控制存储成本。
- **搜索相关性测试**——在调优分词器和打分时，使用 ES Ranking Evaluation API 运行自动化的相关性基准测试。
- **可观测性**——把 p50/p95/p99 查询延迟、零结果率和建议采纳率作为关键的搜索健康指标进行跟踪。

---

## 参考资料

- [Elasticsearch 官方文档](https://www.elastic.co/guide/en/elasticsearch/reference/current/getting-started.html)
- [Debezium 文档](https://debezium.io/documentation/)
- [Canal GitHub 仓库](https://github.com/alibaba/canal)
- [AWS Database Migration Service — CDC](https://docs.aws.amazon.com/dms/latest/userguide/CHAP_Task.CDC.html)
- [Amazon OpenSearch Service](https://aws.amazon.com/opensearch-service/)

---

## 常见问题

### 如何用 Elasticsearch 构建一个全站搜索引擎？

设计一条数据同步管道（通过 binlog 的 CDC 实现实时更新），使用合适的分词器定义索引 mapping，构建带权限过滤的搜索 API，再加上联想输入、热词等功能。

### 什么是 Elasticsearch 的 CDC 同步，为什么要用它？

CDC（变更数据捕获）通过解析 MySQL binlog 检测行级变更，并近实时地同步到 Elasticsearch，让搜索结果保持新鲜，同时避免轮询或双写带来的复杂性。

### 如何在 Elasticsearch 中实现联想输入建议？

使用 Completion Suggester，配合一个类型为 completion 的专用 suggest 字段。它基于内存中的 FST 数据结构做前缀匹配，能以亚毫秒级延迟返回结果。

### Elasticsearch 与 Algolia：该如何选择？

Elasticsearch 提供完全的掌控力，能处理复杂查询并扩展到数十亿文档，但需要运维专业能力。Algolia 是托管服务，开箱即用、体验极佳，但规模上来后成本更高，可定制性也更弱。


---

## 参考资料

- [Elasticsearch reference](https://www.elastic.co/guide/en/elasticsearch/reference/current/index.html) — Elastic
- [IK Analysis plugin for Chinese](https://github.com/infinilabs/analysis-ik) — GitHub
- [Elasticsearch text analyzers](https://www.elastic.co/guide/en/elasticsearch/reference/8.18/analysis-analyzers.html) — Elastic
