Elasticsearch 从 0 到 1(五):MySQL 数据同步与一致性
前面已经能建索引、写 DSL、接入 Spring Boot 了。
但真实项目里还有一个关键问题:
MySQL 里的数据变化后,ES 里的搜索文档怎么同步?
这篇专门讲 MySQL 和 ES 的数据同步。
先看前文:
MySQL 和 ES 谁是主数据
通常情况下:
MySQL 是主数据源ES 是搜索数据源也就是说,商品的真实数据仍然以 MySQL 为准。
ES 里只是存一份适合搜索的冗余数据。
比如 MySQL 里有:
productbrandcategoryinventoryES 里可能是一条搜索文档:
{ "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:
@Transactionalpublic 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。
消息表示例字段:
idbusiness_idevent_typestatusretry_counterror_messagecreated_atupdated_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。
下一篇讲性能和生产环境注意事项:
参考
If this article helped you, please share it with others!
Some information may be outdated






