Skip to content

Elasticsearch 同步与重建索引

商业系统里,ES 通常不是事实主库,而是搜索视图。真正的订单、商品、资产、工单、日志元数据一般保存在 MySQL、PostgreSQL、MongoDB 或业务服务里,然后同步到 ES。

这带来两个核心问题:

  1. 主库数据怎么可靠同步到 ES。
  2. Mapping、分词、字段结构改了以后,怎么平滑重建索引。

学习目标

目标需要掌握什么
知道同步方案同步双写、MQ、Binlog CDC、定时补偿的区别
知道一致性边界ES 通常最终一致,不负责交易强一致
知道怎么防丢可靠消息、重试、死信、补偿、幂等写入
知道怎么防乱序版本号、更新时间、主库快照、按业务 key 顺序
知道怎么重建新建索引、全量导入、增量同步、校验、别名切换
会写 Demo会写同步消费者、upsert、delete、别名切换命令

为什么不能把 ES 当主数据源

ES 的强项是搜索,不是事务事实。

问题如果把 ES 当主库会怎样
写入后近实时可见用户可能短时间搜不到刚写入的数据
更新本质是重写文档高频状态更新成本高
多文档事务弱订单、库存、支付状态不适合放 ES 决策
Mapping 变更成本高字段类型错了通常要重建
依赖搜索副本同步失败会和主库不一致

正确架构:

mermaid
flowchart TD
    A["业务操作"] --> B["主库事务成功"]
    B --> C["产生变更事件"]
    C --> D["同步服务写 ES"]
    D --> E["搜索接口查询 ES"]
    F["交易校验"] --> B

搜索页可以读 ES,真正的交易判断要回主库或业务服务。

常见同步方案对比

方案流程优点风险适合场景
同步双写业务代码写主库后再写 ES简单直观ES 抖动影响主流程,双写失败难补小系统、低风险后台
MQ 异步主库事务后发消息,消费者写 ES解耦、吞吐好、可重试有延迟,需要处理丢失和乱序商品、订单、工单搜索
Binlog CDC监听数据库 Binlog 同步业务侵入小字段聚合、跨表组装复杂数据平台、搜索视图同步
定时补偿定时扫描主库修复 ES兜底可靠延迟高,不能实时与 MQ/CDC 搭配
全量重建从主库全量导入新索引结构变更可靠耗时、资源消耗大Mapping 或分词变更

商业项目通常选择:MQ 或 CDC 做主链路,定时补偿做兜底,重建索引用别名切换。

如果要系统理解 MySQL 和 ES 为什么只能最终一致、双写为什么危险、消息如何防丢、乱序如何处理、删除如何同步、补偿和对账怎么落地,继续看:MySQL 与 ES 数据一致性

MQ 同步流程

mermaid
flowchart TD
    A["商品服务修改商品"] --> B["写 MySQL"]
    B --> C["事务提交"]
    C --> D["发送商品变更消息"]
    D --> E["搜索同步服务消费"]
    E --> F["按商品ID查主库最新快照"]
    F --> G{"商品是否可搜索"}
    G -- "是" --> H["组装搜索文档"]
    H --> I["upsert 到 ES"]
    G -- "否" --> J["删除或标记 ES 文档"]
    I --> K["记录同步成功"]
    J --> K
    E --> L["失败重试 / 死信 / 补偿"]

为什么消息里通常只放 ID,而不是放完整商品数据?

  1. 商品搜索文档往往来自多张表:商品、品牌、类目、库存、销量、标签。
  2. 消息可能乱序,消费时查主库最新快照能减少旧消息覆盖新数据。
  3. 消息体更小,减少 MQ 存储和网络压力。
  4. 业务字段变更后,同步服务可以重新组装文档。

同步消息设计

json
{
  "eventId": "evt-20260704-0001",
  "eventType": "PRODUCT_CHANGED",
  "bizId": "10001",
  "bizType": "PRODUCT",
  "version": 1720000001000,
  "changedAt": "2026-07-04T10:00:01"
}
字段作用
eventId幂等去重
eventType区分新增、修改、删除、上下架
bizId查询主库最新数据
version防乱序覆盖
changedAt排查同步延迟

幂等写入

MQ 消费天然可能重复。同步 ES 必须幂等。

常见做法:

  1. ES 文档 _id 使用业务唯一 ID。
  2. 写入使用 upsert,相同 ID 覆盖同一文档。
  3. 删除和下架使用业务状态判断。
  4. 记录同步日志,便于排查。
  5. 必要时用版本号避免旧消息覆盖新消息。
json
POST /product_search/_update/10001
{
  "doc": {
    "id": 10001,
    "productName": "无线降噪耳机 Pro",
    "stockStatus": "IN_STOCK",
    "updatedAt": "2026-07-04T10:00:01"
  },
  "doc_as_upsert": true
}

乱序问题怎么处理

假设商品连续发生两次变化:

text
10:00:01 商品改价 299 -> 279
10:00:02 商品下架

如果下架消息先消费,改价消息后消费,就可能把已下架商品重新写回 ES。

处理思路:

mermaid
flowchart TD
    A["消费变更消息"] --> B["按 ID 查主库最新快照"]
    B --> C{"主库当前是否可搜索"}
    C -- "是" --> D["写入最新快照"]
    C -- "否" --> E["删除或标记不可搜"]
    D --> F["写入同步版本"]
    E --> F

更稳的做法是消费端永远以主库当前状态为准,而不是完全相信消息里的旧字段。

如果无法每次查主库,也可以用版本号:

json
POST /product_search/_update/10001
{
  "script": {
    "source": """
      if (ctx._source.version == null || params.version >= ctx._source.version) {
        ctx._source.productName = params.productName;
        ctx._source.version = params.version;
      } else {
        ctx.op = 'noop';
      }
    """,
    "params": {
      "productName": "无线降噪耳机 Pro",
      "version": 1720000002000
    }
  },
  "upsert": {
    "productName": "无线降噪耳机 Pro",
    "version": 1720000002000
  }
}

脚本更新会有额外开销,要根据吞吐和一致性要求取舍。

删除和下架怎么同步

删除有两种策略。

策略做法适合场景
物理删除 ES 文档DELETE /index/_doc/id真删除、无保留需求
逻辑标记不可搜字段 searchVisible=false需要审计、恢复、延迟删除

商品下架通常不一定要物理删除,也可以写成:

json
POST /product_search/_update/10001
{
  "doc": {
    "searchVisible": false,
    "updatedAt": "2026-07-04T10:00:02"
  }
}

搜索时加过滤:

json
{
  "term": {
    "searchVisible": true
  }
}

如果业务要求下架后完全不可搜索,物理删除更直接;如果要保留运营或审计字段,逻辑删除更方便。

定时补偿

任何异步链路都可能失败:消息丢失、消费者异常、ES 写入失败、网络超时、字段组装错误。

补偿任务用于定期对齐主库和 ES。

mermaid
flowchart TD
    A["扫描主库最近变更数据"] --> B["查询 ES 对应文档"]
    B --> C{"是否一致"}
    C -- "一致" --> D["跳过"]
    C -- "不一致" --> E["重新组装文档"]
    E --> F["写入或删除 ES"]
    F --> G["记录补偿结果"]

常见补偿维度:

  1. updated_at 扫描最近变更。
  2. 按业务 ID 抽样校验。
  3. 对失败消息死信队列重放。
  4. 对重要字段做主库和 ES 对比。
  5. 对长时间不同步的数据告警。

为什么 Mapping 改错通常要重建索引

字段类型决定底层索引结构。已经写入为 text 的字段,不能直接改成 keyword;已经使用旧 analyzer 建立倒排索引的文本,也不会因为你改 analyzer 自动重分词。

常见需要重建的情况:

变更是否通常要重建
textkeyword
keywordlong
修改 analyzer
新增字段通常不要
新增 keyword 子字段并希望历史数据可用要重建历史数据
调整副本数不需要
调整 refresh_interval不需要

零停机重建索引流程

mermaid
flowchart TD
    A["当前别名 product_search -> v1"] --> B["创建 product_search_v2"]
    B --> C["从主库全量导入 v2"]
    C --> D["增量变更双写或回放到 v2"]
    D --> E["校验数量和抽样结果"]
    E --> F{"校验是否通过"}
    F -- "否" --> G["修复 Mapping、分词或同步逻辑"]
    G --> C
    F -- "是" --> H["原子切换别名到 v2"]
    H --> I["观察搜索指标"]
    I --> J["保留 v1 一段时间后删除"]

关键原则:

  1. 业务代码访问别名,不访问版本索引。
  2. 新索引先全量导入。
  3. 导入期间的增量变更不能丢。
  4. 切换前要校验数量、抽样文档和关键查询结果。
  5. 切换后要能快速回滚到旧索引。

别名切换命令

创建 v1 并绑定别名:

json
PUT /product_search_v1
{
  "mappings": {
    "properties": {
      "id": { "type": "long" },
      "productName": { "type": "text" }
    }
  }
}
json
POST /_aliases
{
  "actions": [
    { "add": { "index": "product_search_v1", "alias": "product_search" } }
  ]
}

切到 v2:

json
POST /_aliases
{
  "actions": [
    { "remove": { "index": "product_search_v1", "alias": "product_search" } },
    { "add": { "index": "product_search_v2", "alias": "product_search" } }
  ]
}

这个操作是原子的。业务请求要么看到旧索引,要么看到新索引,不应该看到中间状态。

全量导入优化

全量导入时更关注吞吐,可以临时调整:

json
PUT /product_search_v2/_settings
{
  "index": {
    "refresh_interval": "-1",
    "number_of_replicas": 0
  }
}

导入完成后恢复:

json
PUT /product_search_v2/_settings
{
  "index": {
    "refresh_interval": "1s",
    "number_of_replicas": 1
  }
}

然后手动 refresh:

json
POST /product_search_v2/_refresh

注意:是否能临时把副本设为 0,要看业务可用性要求。如果导入过程也需要高可用,就不能简单降低副本。

Java 同步消费者 Demo

下面示例表达同步消费者的核心结构:消费消息、查主库、组装文档、幂等写入 ES。

java
public class ProductSearchSyncConsumer {

    private final ProductRepository productRepository;
    private final ProductSearchWriter productSearchWriter;

    public ProductSearchSyncConsumer(ProductRepository productRepository,
                                     ProductSearchWriter productSearchWriter) {
        this.productRepository = productRepository;
        this.productSearchWriter = productSearchWriter;
    }

    public void onMessage(ProductChangedEvent event) {
        ProductSnapshot snapshot = productRepository.findSearchSnapshot(event.productId());

        if (snapshot == null || !snapshot.searchVisible()) {
            productSearchWriter.delete(event.productId());
            return;
        }

        ProductSearchDoc doc = ProductSearchDoc.from(snapshot);
        productSearchWriter.upsert(doc);
    }
}

写 ES 的动作要有重试、失败日志和死信。不要吞异常,否则同步失败后很难发现。

校验怎么做

校验项方法
数量主库可搜索数据量与 ES 文档数量对比
抽样随机抽业务 ID,对比关键字段
查询结果选典型关键词,对比 v1、v2 命中和排序
聚合结果对品牌、类目、状态聚合结果做抽样
延迟统计主库变更到 ES 可见的时间
错误统计同步失败、重试、死信

如果只是数量一致,不代表搜索质量一致。分词、字段权重、过滤条件、排序都可能影响结果。

常见问题

问题原因处理
主库有数据,ES 搜不到同步失败、refresh 未完成、过滤条件误伤查同步日志、ES 文档、DSL
ES 还是旧价格MQ 堆积、旧消息覆盖、补偿未跑查消息延迟、版本号、updatedAt
下架商品还能搜到删除消息失败或逻辑过滤缺失查主库状态和 ES 文档状态
重建后搜索结果变差analyzer、boost、Mapping 与旧版本不同对比 _analyze_explain
切别名后报错新索引字段缺失或别名没切完整回滚别名,补齐 Mapping
全量导入很慢bulk 不合理、refresh 频繁、副本多、DB 慢调整批量、并发、refresh、副本

关联知识点

知识点继续学习什么
底层原理refresh、segment、更新、查询流程
Mapping 与查询为什么字段类型改错要重建
消息堆积与背压ES 同步消息堆积怎么排查
性能优化与排查同步慢、写入慢、查询慢如何定位
商业搜索实践商品、订单、工单、日志场景设计
MySQL 与 ES 数据一致性双写、MQ、CDC、幂等、乱序、补偿、对账

小结

ES 同步的核心是:主库保存事实,ES 保存搜索视图。同步链路要能处理重复、乱序、失败和补偿。Mapping 或分词变更时,不要在旧索引上硬改,而是创建新索引、全量导入、增量追平、校验结果、原子切换别名。这样才能在商业系统里既保证搜索能力,又不破坏主业务稳定性。