Elasticsearch 同步与重建索引
商业系统里,ES 通常不是事实主库,而是搜索视图。真正的订单、商品、资产、工单、日志元数据一般保存在 MySQL、PostgreSQL、MongoDB 或业务服务里,然后同步到 ES。
这带来两个核心问题:
- 主库数据怎么可靠同步到 ES。
- Mapping、分词、字段结构改了以后,怎么平滑重建索引。
学习目标
| 目标 | 需要掌握什么 |
|---|---|
| 知道同步方案 | 同步双写、MQ、Binlog CDC、定时补偿的区别 |
| 知道一致性边界 | ES 通常最终一致,不负责交易强一致 |
| 知道怎么防丢 | 可靠消息、重试、死信、补偿、幂等写入 |
| 知道怎么防乱序 | 版本号、更新时间、主库快照、按业务 key 顺序 |
| 知道怎么重建 | 新建索引、全量导入、增量同步、校验、别名切换 |
| 会写 Demo | 会写同步消费者、upsert、delete、别名切换命令 |
为什么不能把 ES 当主数据源
ES 的强项是搜索,不是事务事实。
| 问题 | 如果把 ES 当主库会怎样 |
|---|---|
| 写入后近实时可见 | 用户可能短时间搜不到刚写入的数据 |
| 更新本质是重写文档 | 高频状态更新成本高 |
| 多文档事务弱 | 订单、库存、支付状态不适合放 ES 决策 |
| Mapping 变更成本高 | 字段类型错了通常要重建 |
| 依赖搜索副本 | 同步失败会和主库不一致 |
正确架构:
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 同步流程
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,而不是放完整商品数据?
- 商品搜索文档往往来自多张表:商品、品牌、类目、库存、销量、标签。
- 消息可能乱序,消费时查主库最新快照能减少旧消息覆盖新数据。
- 消息体更小,减少 MQ 存储和网络压力。
- 业务字段变更后,同步服务可以重新组装文档。
同步消息设计
{
"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 必须幂等。
常见做法:
- ES 文档
_id使用业务唯一 ID。 - 写入使用 upsert,相同 ID 覆盖同一文档。
- 删除和下架使用业务状态判断。
- 记录同步日志,便于排查。
- 必要时用版本号避免旧消息覆盖新消息。
POST /product_search/_update/10001
{
"doc": {
"id": 10001,
"productName": "无线降噪耳机 Pro",
"stockStatus": "IN_STOCK",
"updatedAt": "2026-07-04T10:00:01"
},
"doc_as_upsert": true
}乱序问题怎么处理
假设商品连续发生两次变化:
10:00:01 商品改价 299 -> 279
10:00:02 商品下架如果下架消息先消费,改价消息后消费,就可能把已下架商品重新写回 ES。
处理思路:
flowchart TD
A["消费变更消息"] --> B["按 ID 查主库最新快照"]
B --> C{"主库当前是否可搜索"}
C -- "是" --> D["写入最新快照"]
C -- "否" --> E["删除或标记不可搜"]
D --> F["写入同步版本"]
E --> F更稳的做法是消费端永远以主库当前状态为准,而不是完全相信消息里的旧字段。
如果无法每次查主库,也可以用版本号:
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 | 需要审计、恢复、延迟删除 |
商品下架通常不一定要物理删除,也可以写成:
POST /product_search/_update/10001
{
"doc": {
"searchVisible": false,
"updatedAt": "2026-07-04T10:00:02"
}
}搜索时加过滤:
{
"term": {
"searchVisible": true
}
}如果业务要求下架后完全不可搜索,物理删除更直接;如果要保留运营或审计字段,逻辑删除更方便。
定时补偿
任何异步链路都可能失败:消息丢失、消费者异常、ES 写入失败、网络超时、字段组装错误。
补偿任务用于定期对齐主库和 ES。
flowchart TD
A["扫描主库最近变更数据"] --> B["查询 ES 对应文档"]
B --> C{"是否一致"}
C -- "一致" --> D["跳过"]
C -- "不一致" --> E["重新组装文档"]
E --> F["写入或删除 ES"]
F --> G["记录补偿结果"]常见补偿维度:
- 按
updated_at扫描最近变更。 - 按业务 ID 抽样校验。
- 对失败消息死信队列重放。
- 对重要字段做主库和 ES 对比。
- 对长时间不同步的数据告警。
为什么 Mapping 改错通常要重建索引
字段类型决定底层索引结构。已经写入为 text 的字段,不能直接改成 keyword;已经使用旧 analyzer 建立倒排索引的文本,也不会因为你改 analyzer 自动重分词。
常见需要重建的情况:
| 变更 | 是否通常要重建 |
|---|---|
text 改 keyword | 要 |
keyword 改 long | 要 |
| 修改 analyzer | 要 |
| 新增字段 | 通常不要 |
| 新增 keyword 子字段并希望历史数据可用 | 要重建历史数据 |
| 调整副本数 | 不需要 |
| 调整 refresh_interval | 不需要 |
零停机重建索引流程
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 一段时间后删除"]关键原则:
- 业务代码访问别名,不访问版本索引。
- 新索引先全量导入。
- 导入期间的增量变更不能丢。
- 切换前要校验数量、抽样文档和关键查询结果。
- 切换后要能快速回滚到旧索引。
别名切换命令
创建 v1 并绑定别名:
PUT /product_search_v1
{
"mappings": {
"properties": {
"id": { "type": "long" },
"productName": { "type": "text" }
}
}
}POST /_aliases
{
"actions": [
{ "add": { "index": "product_search_v1", "alias": "product_search" } }
]
}切到 v2:
POST /_aliases
{
"actions": [
{ "remove": { "index": "product_search_v1", "alias": "product_search" } },
{ "add": { "index": "product_search_v2", "alias": "product_search" } }
]
}这个操作是原子的。业务请求要么看到旧索引,要么看到新索引,不应该看到中间状态。
全量导入优化
全量导入时更关注吞吐,可以临时调整:
PUT /product_search_v2/_settings
{
"index": {
"refresh_interval": "-1",
"number_of_replicas": 0
}
}导入完成后恢复:
PUT /product_search_v2/_settings
{
"index": {
"refresh_interval": "1s",
"number_of_replicas": 1
}
}然后手动 refresh:
POST /product_search_v2/_refresh注意:是否能临时把副本设为 0,要看业务可用性要求。如果导入过程也需要高可用,就不能简单降低副本。
Java 同步消费者 Demo
下面示例表达同步消费者的核心结构:消费消息、查主库、组装文档、幂等写入 ES。
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 或分词变更时,不要在旧索引上硬改,而是创建新索引、全量导入、增量追平、校验结果、原子切换别名。这样才能在商业系统里既保证搜索能力,又不破坏主业务稳定性。
