mobile wallpaper 1mobile wallpaper 2mobile wallpaper 3mobile wallpaper 4
1314 words
3 minutes
Elasticsearch 从 0 到 1(五):MySQL 数据同步与一致性
2025-06-19

Elasticsearch 从 0 到 1(五):MySQL 数据同步与一致性#

前面已经能建索引、写 DSL、接入 Spring Boot 了。

但真实项目里还有一个关键问题:

MySQL 里的数据变化后,ES 里的搜索文档怎么同步?

这篇专门讲 MySQL 和 ES 的数据同步。

先看前文:

MySQL 和 ES 谁是主数据#

通常情况下:

MySQL 是主数据源
ES 是搜索数据源

也就是说,商品的真实数据仍然以 MySQL 为准。

ES 里只是存一份适合搜索的冗余数据。

比如 MySQL 里有:

product
brand
category
inventory

ES 里可能是一条搜索文档:

{
"productId": 1001,
"title": "无线蓝牙耳机",
"brand": "SoundGo",
"category": "数码耳机",
"price": 199,
"stock": 80,
"status": "ON_SALE"
}

ES 文档是为了搜索而设计的,不是为了替代 MySQL。

全量同步#

全量同步就是把 MySQL 里的数据全部重新导入 ES。

适合:

  • 第一次初始化 ES;
  • Mapping 变更后重建索引;
  • 数据修复;
  • ES 数据丢失后重建。

流程:

分页查询 MySQL 商品数据
组装 ProductDocument
批量写入 ES
记录同步进度

伪代码:

int pageNo = 1;
int pageSize = 500;
while (true) {
List<Product> products = productMapper.page(pageNo, pageSize);
if (products.isEmpty()) {
break;
}
List<ProductDocument> docs = products.stream()
.map(this::buildDocument)
.toList();
productSearchService.bulkSave(docs);
pageNo++;
}

注意不要一次性把所有商品读进内存。

批量写入 ES#

单条写入效率低,全量同步应该使用 bulk。

示例:

BulkRequest.Builder br = new BulkRequest.Builder();
for (ProductDocument doc : documents) {
br.operations(op -> op.index(idx -> idx
.index("products")
.id(String.valueOf(doc.getProductId()))
.document(doc)
));
}
BulkResponse response = client.bulk(br.build());

写完要检查是否有失败项:

if (response.errors()) {
response.items().forEach(item -> {
if (item.error() != null) {
log.error("bulk index failed, id={}, reason={}",
item.id(), item.error().reason());
}
});
}

不要只看请求没有抛异常就认为全部成功。

增量同步#

增量同步是指商品发生变化时,同步更新 ES。

比如:

  • 商品新增;
  • 商品标题修改;
  • 价格变化;
  • 上下架;
  • 库存变化;
  • 分类或品牌变化。

最简单的方式是在业务代码里同步写 ES:

@Transactional
public void updateProduct(ProductUpdateRequest request) {
productMapper.update(request);
productSearchService.save(buildDocument(request.getProductId()));
}

但这种方式有问题:

  • MySQL 更新成功,ES 更新失败怎么办;
  • ES 更新成功,事务回滚怎么办;
  • ES 慢了会拖慢主流程;
  • 重试和补偿不好做。

所以核心业务里不建议简单粗暴地同步双写。

方案一:业务代码 + MQ#

更常见的方式是:

更新 MySQL
发送商品变更消息
消费者查询 MySQL 最新数据
更新 ES

比如商品更新后发消息:

productEventPublisher.publish(new ProductChangedEvent(productId));

消费者:

public void consume(ProductChangedEvent event) {
Product product = productMapper.selectById(event.productId());
ProductDocument document = buildDocument(product);
productSearchService.save(document);
}

消费者最好查询 MySQL 最新数据,而不是完全相信消息里的旧数据。

方案二:本地消息表#

如果担心“数据库提交成功,但 MQ 发送失败”,可以用本地消息表。

一个事务里同时写:

1. 更新商品表
2. 写 product_sync_message 表

事务提交后,由后台任务扫描消息表,发送到 MQ 或直接同步 ES。

消息表示例字段:

id
business_id
event_type
status
retry_count
error_message
created_at
updated_at

这样即使同步失败,也可以重试。

这和异步可靠性是一类问题,可以参考:Spring @Async 实战:异步方法、线程池配置与常见失效场景

方案三:监听 Binlog#

还有一种方式是监听 MySQL Binlog。

常见工具是 Canal。

流程:

MySQL Binlog
Canal 解析变更
消息队列
同步服务更新 ES

优点:

  • 对业务代码侵入小;
  • 可以捕获数据库层面的变更;
  • 适合多个系统共享同步链路。

缺点:

  • 架构更复杂;
  • 运维成本更高;
  • 需要处理表结构变更;
  • 需要把数据库变更转换成 ES 文档。

小项目可以先用业务事件,大项目再考虑 Binlog 方案。

删除和下架怎么同步#

商品下架有两种处理方式。

第一种:删除 ES 文档。

client.delete(d -> d
.index("products")
.id(String.valueOf(productId))
);

第二种:保留文档,但改状态。

{
"status": "OFF_SALE"
}

搜索时过滤:

{ "term": { "status": "ON_SALE" } }

我更倾向第二种,因为后台可能还需要查到下架商品,而且状态变更比删除更容易追踪。

一致性怎么理解#

MySQL 和 ES 通常不是强一致,而是最终一致。

也就是说:

MySQL 更新成功后
ES 可能延迟几百毫秒或几秒才更新

大多数搜索场景可以接受这种延迟。

但要注意:

  • 下架商品不能继续被用户买到;
  • 价格变更不能长期不同步;
  • 库存最好不要完全依赖 ES 判断;
  • 支付、订单、库存扣减不要以 ES 为准。

ES 更适合做搜索,不适合做核心交易判断。

同步失败怎么处理#

同步失败要能发现、能重试、能补偿。

建议至少有:

  • 失败日志;
  • 重试次数;
  • 死信队列或失败表;
  • 手动补偿接口;
  • 定时校验任务;
  • 同步延迟监控。

比如定时校验:

每天凌晨抽样对比 MySQL 和 ES
发现缺失或字段不一致
重新同步

不要让 ES 不一致问题长期静默存在。

这一篇先记住什么#

  • MySQL 通常是主数据源,ES 是搜索冗余;
  • 全量同步适合初始化、重建索引、数据修复;
  • 增量同步可以用业务事件、MQ、本地消息表、Binlog;
  • 不建议核心流程简单同步双写 MySQL 和 ES;
  • ES 和 MySQL 通常是最终一致;
  • 同步失败要有重试、告警和补偿;
  • 搜索可以依赖 ES,交易判断不要依赖 ES。

下一篇讲性能和生产环境注意事项:

Elasticsearch 从 0 到 1(六):性能优化与生产环境注意事项

参考#

Share

If this article helped you, please share it with others!

Elasticsearch 从 0 到 1(五):MySQL 数据同步与一致性
https://mizuki.mysqil.com/posts/elasticsearch-05-mysql-sync/
Author
梦幻晨风
Published at
2025-06-19
License
CC BY-NC-SA 4.0

Some information may be outdated

Table of Contents