MySQL 与 ES 数据一致性
商业系统里,MySQL 通常是事实数据源,Elasticsearch 是搜索视图。事实数据源负责事务、约束、状态机和最终正确性;搜索视图负责全文检索、聚合筛选、相关性排序和高效查询。
一句话理解:
MySQL 和 ES 很难做强一致,商业项目通常做最终一致:主库先正确提交,再通过可靠同步、幂等消费、乱序控制、补偿任务、对账校验,把 ES 修正到和主库一致。
学习目标
| 目标 | 需要掌握什么 |
|---|---|
| 知道定位 | MySQL 是事实源,ES 是搜索副本 |
| 知道风险 | 双写失败、消息丢失、重复消费、乱序、删除漏同步、refresh 延迟 |
| 知道方案 | 同步双写、MQ、Binlog CDC、本地消息表、事务消息、定时补偿 |
| 知道原理 | 为什么最终一致,为什么消费端要查主库最新快照 |
| 会写 Demo | 表设计、事件消息、消费者 upsert/delete、补偿任务 |
| 会排查 | 搜不到、搜到旧数据、下架仍可搜、同步延迟、消息堆积 |
如果你想把这一页的同步、一致性、失败补偿和重建流程放到完整商业搜索链路里练习,继续做:Elasticsearch 商业场景训练营。
为什么不能要求强一致
MySQL 和 ES 是两个不同系统。一次业务写入如果同时写 MySQL 和 ES,就会遇到典型分布式一致性问题:
flowchart TD
A["业务修改商品"] --> B["写 MySQL"]
B --> C["写 ES"]
B --> D["MySQL 成功"]
C --> E["ES 可能成功、失败或超时"]如果 MySQL 成功但 ES 失败,就不一致;如果 ES 成功但 MySQL 事务回滚,也不一致。除非引入强一致分布式事务,但 ES 本身并不是常规 XA 事务资源,搜索场景也通常不值得为了“立即一致”牺牲主流程可用性。
所以商业系统一般按数据风险分类:
| 数据类型 | 一致性要求 | 推荐读源 |
|---|---|---|
| 支付、库存、账户、订单最终状态 | 强一致或强约束 | MySQL 或业务服务 |
| 商品搜索、订单检索、资产搜索 | 最终一致 | ES |
| 报表、看板、运营统计 | 最终一致 | ES、数仓、汇总表 |
| 用户刚提交后的结果确认 | 尽量读己之写 | MySQL 或同步状态 |
ES 搜索结果可以短暂落后,但不能长期错误;核心业务决策不能依赖 ES。
一致性边界
MySQL 到 ES 的一致性通常追求这几个目标:
| 目标 | 说明 |
|---|---|
| 不丢变更 | 主库有变更,最终要同步到 ES |
| 可重试 | 写 ES 失败可以自动重试 |
| 幂等 | 同一事件重复消费不会写错 |
| 防乱序 | 旧事件不能覆盖新状态 |
| 可补偿 | 已知失败和未知漏同步都能修复 |
| 可观测 | 能看到延迟、失败、积压、差异 |
不要承诺:
- MySQL 提交后 ES 立即可搜索。
- ES 搜索结果永远等于 MySQL 当前状态。
- 消息只投递一次。
- 消费顺序永远和业务写入顺序一致。
常见同步方案对比
| 方案 | 流程 | 优点 | 风险 | 推荐程度 |
|---|---|---|---|---|
| 同步双写 | 写 MySQL 后立即写 ES | 简单 | 主流程受 ES 影响,失败难补 | 低 |
| MQ 异步 | 写库后发消息,消费者写 ES | 解耦、可重试 | 有延迟,要防丢和乱序 | 高 |
| 本地消息表 | 业务数据和消息表同事务提交 | 解决写库成功但消息没发 | 多一张表和投递任务 | 高 |
| RocketMQ 事务消息 | 半消息 + 本地事务 + 提交消息 | 减少消息与事务不一致 | 依赖 MQ 能力和回查 | 中高 |
| Binlog CDC | 监听 MySQL binlog 同步 | 业务侵入小 | 跨表组装复杂,链路复杂 | 高 |
| 定时补偿 | 扫主库对比 ES 并修复 | 兜底可靠 | 延迟高 | 必须有 |
| 全量重建 | 从主库重建新索引 | 修复结构和历史差异 | 成本高 | 结构变更时用 |
成熟方案通常是:
flowchart TD
A["MySQL 本地事务"] --> B["变更事件"]
B --> C["MQ 或 CDC"]
C --> D["同步服务写 ES"]
D --> E["失败重试和死信"]
F["定时补偿"] --> D
G["对账校验"] --> F一致性链路要分别保证什么
MySQL 到 ES 的同步不是“把消息发出去”这么简单,而是一条端到端链路。任何一段没有兜底,最后都会表现成“搜索页看到旧数据”。
flowchart TD
A["业务写入"] --> B["事件产生"]
B --> C["事件投递"]
C --> D["事件消费"]
D --> E["写入 ES"]
E --> F["搜索可见"]
F --> G["补偿对账"]每一段要解决的问题不同:
| 链路 | 要保证什么 | 常见手段 | 如果不处理会怎样 |
|---|---|---|---|
| 业务写入 | MySQL 数据先正确提交 | 本地事务、唯一约束、状态机校验 | 主数据本身就是错的,ES 同步再准也没意义 |
| 事件产生 | 主库变更一定有同步事件 | 本地消息表、事务消息、CDC | MySQL 成功但没有消息,ES 永远旧 |
| 事件投递 | 事件能送到同步服务 | MQ 重试、确认机制、死信、投递日志 | 消息丢失或长期堆积 |
| 事件消费 | 重复和乱序不会写错 | 幂等表、业务主键、版本号、按 key 顺序 | 重复消费报错,旧事件覆盖新状态 |
| 写入 ES | 写失败可以重试 | Bulk、失败明细、退避重试、死信重放 | 单条失败混在批次里没人发现 |
| 搜索可见 | 写入后理解 refresh 延迟 | refresh 窗口、refresh=wait_for、读 MySQL 兜底 | 误把近实时延迟当成同步失败 |
| 补偿对账 | 未知差异能被发现和修复 | 定时扫描、抽样对账、版本对比、人工重放 | 小概率错误沉淀成长期脏数据 |
所以“数据一致性”不是一个 API,而是一套工程机制:可靠事件、幂等消费、乱序控制、失败重试、补偿对账、读路径兜底。
为什么同步双写不推荐
同步双写看起来最简单:
public void updateProduct(Product product) {
productMapper.update(product);
elasticsearchClient.update(product);
}但它有很多失败窗口:
flowchart TD
A["开始修改商品"] --> B["MySQL 更新成功"]
B --> C{"写 ES 是否成功"}
C -- "成功" --> D["暂时一致"]
C -- "失败" --> E["MySQL 新,ES 旧"]
C -- "超时" --> F["不知道 ES 是否成功"]同步双写的问题:
- ES 抖动会拖慢主业务接口。
- MySQL 成功、ES 失败后需要补偿。
- ES 超时时不知道到底写没写成功。
- 如果业务事务后来回滚,ES 可能已经写入错误数据。
- 高并发更新时容易发生旧请求覆盖新请求。
如果系统很小、后台低频、允许人工修复,可以临时用同步双写;核心商业链路不建议。
失败窗口怎么分析
设计一致性方案时,不要只看“正常流程”,要按失败窗口推演。
| 失败点 | 可能结果 | 解决办法 |
|---|---|---|
| MySQL 事务提交前服务宕机 | 业务数据没提交,事件也不应产生 | 本地事务保证一起回滚 |
| MySQL 提交成功但还没发 MQ 就宕机 | 主库新,ES 旧 | 本地消息表保存待发送事件,后台继续投递 |
| MQ 已收到但业务事务回滚 | ES 可能同步了不存在或错误状态 | 先提交事务再投递,或使用事务消息回查 |
| MQ 重复投递 | 同一事件被消费多次 | event_id 唯一约束和 ES 业务 _id 幂等 |
| 消息乱序 | 旧价格覆盖新价格,下架后又可搜索 | 消费端查主库最新快照,或使用 version 防旧覆盖 |
| ES Bulk 部分失败 | 一批里有的成功有的失败 | 逐条检查 Bulk item,失败明细入库重试 |
| 消费者一直失败 | 消息堆积,ES 长期旧 | 重试退避、死信队列、告警、人工重放 |
| 全量重建期间有增量 | 新索引切换后缺最新数据 | 记录起点,增量追平,校验后切别名 |
能把这些失败窗口讲清楚,才算真正理解 MySQL 和 ES 的一致性。
RocketMQ 事务消息怎么用在同步里
RocketMQ 事务消息可以降低“本地事务和消息发送不一致”的风险。它不是让 MySQL 和 ES 强一致,而是尽量保证“业务成功后消息最终可见”。
flowchart TD
A["发送半消息"] --> B["执行 MySQL 本地事务"]
B --> C{"本地事务结果"}
C -- "提交成功" --> D["提交消息"]
C -- "回滚失败" --> E["回滚消息"]
C -- "未知" --> F["Broker 回查事务状态"]
F --> G["根据本地事务记录提交或回滚"]
D --> H["同步服务消费并写 ES"]关键点:
| 点 | 说明 |
|---|---|
| 半消息 | Broker 先保存消息,但消费者暂时不可见 |
| 本地事务 | 业务服务写 MySQL,并记录可回查的事务状态 |
| 提交消息 | 本地事务成功后,消息变成可消费 |
| 回滚消息 | 本地事务失败后,消息丢弃 |
| 事务回查 | 生产者宕机或网络异常时,Broker 回查本地事务状态 |
事务消息仍然要求消费端幂等。因为消息可能重复投递,消费者写 ES 也可能失败重试。它解决的是“消息产生可靠性”,不是“一次消费绝对成功”。
如果团队已经有统一本地消息表,继续用本地消息表也可以;如果消息中间件能力成熟,并且业务代码能可靠实现事务回查,RocketMQ 事务消息也适合。
推荐架构一:本地消息表加 MQ
本地消息表解决的是“业务提交成功,但消息没发出去”的问题。
flowchart TD
A["业务服务"] --> B["开启 MySQL 事务"]
B --> C["更新业务表"]
C --> D["插入 outbox 事件表"]
D --> E["提交事务"]
E --> F["投递任务发送 MQ"]
F --> G["同步服务消费"]
G --> H["写 ES"]业务数据和事件表在同一个 MySQL 本地事务里提交,要么都成功,要么都失败。
表设计 Demo
商品表:
create table product (
id bigint primary key,
product_name varchar(200) not null,
status varchar(32) not null,
price int not null,
version bigint not null,
updated_at datetime not null
) engine = InnoDB default charset = utf8mb4;事件表:
create table search_sync_event (
id bigint primary key auto_increment,
event_id varchar(64) not null,
biz_type varchar(32) not null,
biz_id varchar(64) not null,
event_type varchar(32) not null,
version bigint not null,
status varchar(32) not null,
retry_count int not null default 0,
next_retry_time datetime not null,
created_at datetime not null,
updated_at datetime not null,
unique key uk_event_id(event_id),
key idx_status_retry(status, next_retry_time),
key idx_biz(biz_type, biz_id)
) engine = InnoDB default charset = utf8mb4;事件字段含义:
| 字段 | 作用 |
|---|---|
event_id | 消息幂等去重 |
biz_type | 商品、订单、资产、工单等类型 |
biz_id | 业务主键 |
event_type | 创建、修改、删除、上下架 |
version | 防乱序覆盖 |
status | 待发送、已发送、失败 |
retry_count | 投递重试次数 |
next_retry_time | 下次重试时间 |
写业务和事件
public void updateProduct(ProductUpdateCommand command) {
Product product = productMapper.selectById(command.getProductId());
long nextVersion = product.getVersion() + 1;
productMapper.updateProduct(
command.getProductId(),
command.getProductName(),
command.getPrice(),
nextVersion
);
SearchSyncEvent event = new SearchSyncEvent();
event.setEventId(UUID.randomUUID().toString());
event.setBizType("PRODUCT");
event.setBizId(String.valueOf(command.getProductId()));
event.setEventType("PRODUCT_CHANGED");
event.setVersion(nextVersion);
event.setStatus("WAIT_SEND");
event.setNextRetryTime(new Date());
searchSyncEventMapper.insert(event);
}这段代码必须放在同一个本地事务里。Spring 项目中可以由 @Transactional 包住业务更新和事件插入。
投递任务
public void publishPendingEvents() {
List<SearchSyncEvent> events = searchSyncEventMapper.findWaitSend(100);
for (SearchSyncEvent event : events) {
try {
mqProducer.send("search.sync.event", event);
searchSyncEventMapper.markSent(event.getId());
} catch (Exception ex) {
searchSyncEventMapper.markSendFailed(event.getId(), nextRetryTime(event));
}
}
}投递任务要支持:
- 批量拉取。
- 失败重试。
- 最大重试次数。
- 告警。
- 人工重投。
推荐架构二:Binlog CDC
CDC Change Data Capture 的思路是监听 MySQL binlog,把数据库变更转换为同步事件。常见工具有 Canal、Debezium、Flink CDC 等。
flowchart TD
A["业务服务写 MySQL"] --> B["MySQL binlog"]
B --> C["CDC 组件"]
C --> D["变更事件"]
D --> E["同步服务"]
E --> F["组装 ES 文档"]
F --> G["写 ES"]CDC 的优点:
- 对业务代码侵入小。
- 不需要每个业务方法手动发消息。
- 能捕获直接 SQL 修改带来的变更。
- 适合数据平台、搜索平台统一同步。
CDC 的难点:
| 难点 | 说明 |
|---|---|
| 跨表组装 | ES 文档常来自多张表,单表 binlog 不一定够 |
| 变更顺序 | 要按主键或业务维度处理乱序 |
| 删除事件 | 要正确处理 delete 和逻辑删除 |
| 全量加增量 | 第一次初始化要处理全量期间的增量 |
| 位点保存 | CDC 消费位点要可靠保存 |
| 字段语义 | binlog 只有字段变化,不一定知道业务意图 |
CDC 适合“统一数据同步平台”,业务消息适合“业务语义清楚的事件驱动”。很多商业项目会二者结合。
消费端为什么要查 MySQL 最新快照
同步消息有两种设计:
| 方式 | 内容 | 风险 |
|---|---|---|
| 消息放完整文档 | 商品名、价格、库存、标签全放消息里 | 消息大,字段变化难维护,乱序风险高 |
| 消息只放 ID 和版本 | 消费时查 MySQL 最新状态 | 多一次查库,但更稳 |
推荐消费端按 ID 查询 MySQL 最新快照:
flowchart TD
A["消费商品变更事件"] --> B["按商品 ID 查 MySQL 最新快照"]
B --> C{"当前是否可搜索"}
C -- "是" --> D["组装最新 ES 文档"]
D --> E["upsert 到 ES"]
C -- "否" --> F["删除或标记不可搜"]这样可以减少乱序问题。比如商品先改价,再下架。如果“改价消息”晚于“下架消息”被消费,消费端查 MySQL 最新状态发现商品已经下架,就不会把旧价格文档重新写回可搜索状态。
幂等设计
MQ 和 CDC 都可能重复投递。消费端必须做到重复执行结果一致。
文档 ID 幂等
ES 文档 _id 使用业务主键:
PUT /product_search/_doc/10001
{
"id": 10001,
"productName": "无线蓝牙耳机",
"version": 101
}同一个商品重复写入同一个 _id,不会生成多份文档。
消费日志幂等
同步日志表:
create table search_sync_log (
id bigint primary key auto_increment,
event_id varchar(64) not null,
biz_type varchar(32) not null,
biz_id varchar(64) not null,
version bigint not null,
status varchar(32) not null,
error_message varchar(500),
created_at datetime not null,
unique key uk_event_id(event_id),
key idx_biz(biz_type, biz_id)
) engine = InnoDB default charset = utf8mb4;消费前先插入 event_id,如果唯一约束冲突,说明已经处理过或正在处理。
乱序控制
乱序是 ES 同步最容易被问到的问题。
典型场景:
10:00:01 商品改价 version=101
10:00:02 商品下架 version=102如果 version=102 先写 ES,version=101 后写 ES,就可能旧数据覆盖新数据。
解决方式:
| 方式 | 思路 | 适合场景 |
|---|---|---|
| 消费时查主库最新快照 | 永远以 MySQL 当前状态组装文档 | 最常用 |
| 业务版本号 | 旧版本不能覆盖新版本 | 高频更新 |
| 按业务 key 顺序消费 | 同一个商品进入同一分区 | Kafka、RocketMQ 顺序场景 |
| ES 外部版本 | 用版本控制写入 | 对版本管理清晰的场景 |
版本号脚本 Demo
POST /product_search/_update/10001
{
"script": {
"source": """
if (ctx._source.version == null || params.version >= ctx._source.version) {
ctx._source.productName = params.productName;
ctx._source.status = params.status;
ctx._source.version = params.version;
ctx._source.updatedAt = params.updatedAt;
} else {
ctx.op = 'noop';
}
""",
"params": {
"productName": "无线蓝牙耳机",
"status": "ON_SALE",
"version": 102,
"updatedAt": "2026-07-04T10:00:02"
}
},
"upsert": {
"productName": "无线蓝牙耳机",
"status": "ON_SALE",
"version": 102,
"updatedAt": "2026-07-04T10:00:02"
}
}脚本更新有额外开销。如果吞吐很高,优先考虑“查主库最新快照 + 按业务 key 顺序消费”。
删除和下架一致性
删除比更新更容易出事故。
常见错误:
- MySQL 已删除,ES 文档没删,搜索还能搜到。
- 商品下架只是改状态,但 ES 查询没有过滤状态。
- 物理删除消息乱序后,旧更新又把文档写回来。
- 全量重建时把已删除数据重新导入。
推荐策略:
| 业务动作 | ES 处理 |
|---|---|
| 商品下架 | 设置 searchVisible=false 或删除文档 |
| 商品软删除 | 设置删除状态,并搜索时过滤 |
| 商品物理删除 | 删除 ES 文档 |
| 订单状态变化 | 更新状态字段,但订单事实仍查 MySQL |
逻辑不可搜:
POST /product_search/_update/10001
{
"doc": {
"searchVisible": false,
"status": "OFF_SALE",
"version": 103,
"updatedAt": "2026-07-04T10:00:03"
}
}搜索 DSL 必须过滤:
{
"term": {
"searchVisible": true
}
}物理删除:
DELETE /product_search/_doc/10001删除也要幂等。重复删除同一个文档应该被视为成功或可忽略。
Refresh 延迟不是同步失败
ES 写入成功后,不代表立刻能被搜索到。只有 refresh 生成可搜索 segment 后,搜索请求才能看到新文档。
flowchart TD
A["写入 ES 成功"] --> B["进入内存 buffer 和 translog"]
B --> C["等待 refresh"]
C --> D["生成可搜索 segment"]
D --> E["搜索可见"]常见误判:
| 现象 | 可能原因 |
|---|---|
| 按 ID GET 能查到,搜索查不到 | refresh 未完成或查询条件不匹配 |
| 写入后 1 秒内搜不到 | 近实时正常窗口 |
| 长时间搜不到 | 同步失败、Mapping、分词、过滤条件问题 |
如果某个后台操作要求“提交后立刻可见”,可以:
- 操作完成页直接读 MySQL。
- 给用户展示“同步中”状态。
- 对少量关键写入使用
refresh=wait_for。 - 不要全局把
refresh_interval调得极短,否则写入和 merge 压力会上升。
读路径怎么设计
搜索列表读 ES,详情和交易校验读 MySQL。
flowchart TD
A["用户搜索"] --> B["查询 ES"]
B --> C["返回命中的业务 ID"]
C --> D{"是否需要强实时详情"}
D -- "需要" --> E["按 ID 回 MySQL 补最新详情"]
D -- "不需要" --> F["直接返回 ES 冗余字段"]常见策略:
| 场景 | 读法 |
|---|---|
| 商品搜索列表 | ES 返回列表字段 |
| 商品详情页 | MySQL 或商品服务返回事实详情 |
| 下单校验 | 库存和价格回业务服务确认 |
| 后台订单检索 | ES 查列表,详情回 MySQL |
| 用户刚修改后立即查看 | MySQL 或同步状态兜底 |
这样做的好处是:ES 可以短暂落后,但不会影响关键业务正确性。
补偿任务
只靠 MQ 或 CDC 不够,必须有补偿。
补偿分两类:
| 类型 | 解决什么 |
|---|---|
| 已知失败补偿 | 消费异常、写 ES 失败、死信队列 |
| 未知差异对账 | 消息漏了、代码 bug、人工改库、历史脏数据 |
按更新时间补偿
flowchart TD
A["定时扫描 MySQL 最近变更"] --> B["查询对应 ES 文档"]
B --> C{"关键字段是否一致"}
C -- "一致" --> D["跳过"]
C -- "不一致" --> E["重新写入 ES"]
E --> F["记录补偿日志"]SQL:
select id, version, updated_at
from product
where updated_at >= ?
and updated_at < ?
order by updated_at asc
limit 1000;补偿要注意:
- 分批扫描,避免压垮主库。
- 按时间窗口推进,保存游标。
- 对失败记录重试。
- 对长期失败告警。
- 不要和全量重建互相冲突。
抽样对账
对账不一定每次全量,可以抽样或按业务范围:
| 对账方式 | 适合场景 |
|---|---|
| 最近 1 小时变更全量对比 | 发现同步延迟和失败 |
| 每天抽样 1 万个 ID | 发现长期差异 |
| 按租户或机构分批对账 | 多租户数据隔离 |
| 按状态对账 | 下架、删除、异常状态重点查 |
| 全量对账 | 大版本发布或重大修复后 |
全量重建时的一致性
Mapping、分词器、字段类型改了,通常要重建索引。重建期间也会有业务增量变更,如果处理不好,新索引切过去就是旧数据。
推荐流程:
flowchart TD
A["创建新索引 v2"] --> B["记录全量开始时间 T1"]
B --> C["从 MySQL 全量导入 v2"]
C --> D["回放 T1 之后的增量"]
D --> E["校验数量和抽样字段"]
E --> F{"校验通过"}
F -- "否" --> G["修复并重新追平"]
F -- "是" --> H["原子切换别名"]关键点:
- 业务访问别名,不写死版本索引。
- 全量导入期间的增量不能丢。
- 切换前必须校验。
- 切换后保留旧索引,方便回滚。
- 如果使用 CDC,要保存 binlog 位点;如果使用 MQ,要保存事件时间或版本。
Java 消费者 Demo
下面示例使用 JDK 8 风格,避免依赖高版本语法。
public class ProductSearchSyncConsumer {
private final ProductRepository productRepository;
private final ProductSearchWriter productSearchWriter;
private final SearchSyncLogRepository syncLogRepository;
public ProductSearchSyncConsumer(ProductRepository productRepository,
ProductSearchWriter productSearchWriter,
SearchSyncLogRepository syncLogRepository) {
this.productRepository = productRepository;
this.productSearchWriter = productSearchWriter;
this.syncLogRepository = syncLogRepository;
}
public void onMessage(SearchSyncEvent event) {
if (!syncLogRepository.tryStart(event.getEventId())) {
return;
}
try {
ProductSnapshot snapshot = productRepository.findSearchSnapshot(event.getBizId());
if (snapshot == null || !snapshot.isSearchVisible()) {
productSearchWriter.delete(event.getBizId());
syncLogRepository.markSuccess(event.getEventId());
return;
}
ProductSearchDoc doc = ProductSearchDoc.from(snapshot);
productSearchWriter.upsert(doc);
syncLogRepository.markSuccess(event.getEventId());
} catch (Exception ex) {
syncLogRepository.markFailed(event.getEventId(), ex.getMessage());
throw ex;
}
}
}这个 Demo 体现几个原则:
- 用
eventId做消费幂等。 - 消费时查 MySQL 最新快照。
- 当前不可搜索就删除或标记 ES 文档。
- 失败不吞异常,让 MQ 重试或进入死信。
- 同步日志落库,便于排查。
ES 更新失败怎么办
ES 更新失败时,第一原则是:不能影响 MySQL 主事务的正确提交,也不能把失败静默吞掉。因为 MySQL 是事实源,ES 是搜索视图。ES 短暂失败可以接受,长期失败必须被发现、重试、补偿和告警。
整体处理流程:
flowchart TD
A["消费同步事件"] --> B["查 MySQL 最新快照"]
B --> C["组装 ES 文档"]
C --> D["写入 ES"]
D --> E{"是否成功"}
E -- "成功" --> F["记录同步成功"]
E -- "可重试失败" --> G["记录失败原因"]
G --> H["抛异常让 MQ 重试"]
H --> I{"超过重试次数"}
I -- "否" --> A
I -- "是" --> J["进入死信队列"]
J --> K["告警和人工重放"]
E -- "不可重试失败" --> L["修正数据或 Mapping"]
L --> M["补偿任务重新同步"]先区分失败类型
不是所有 ES 更新失败都用同一种方式处理。先判断失败类型,才能决定重试、修数据还是重建索引。
| 失败类型 | 常见原因 | 处理方式 |
|---|---|---|
| 网络超时 | ES 短暂抖动、网络抖动、客户端超时太短 | 记录失败,按退避策略重试 |
| 线程池拒绝 | 写入并发太高,ES write 线程池或队列满 | 降低消费并发、调小 Bulk、扩容 ES |
| 集群只读 | 磁盘超过水位,索引被设置只读 | 清理磁盘、扩容、解除只读 |
| Mapping 冲突 | 字段类型和 Mapping 不一致 | 修正数据转换或重建索引 |
| 文档过大 | 单条文档字段过多、内容过长 | 裁剪字段、拆文档、限制冗余 |
| 版本冲突 | 并发更新同一文档,旧版本覆盖新版本 | 使用版本号、脚本更新或查主库最新快照 |
| 认证权限失败 | 账号权限、证书、网络策略变更 | 修复配置,不要盲目重试 |
| Bulk 部分失败 | 一批里部分文档失败 | 逐条检查失败项,只重试失败文档 |
可重试失败通常是网络、超时、429、部分 5xx、临时拒绝;不可重试失败通常是 Mapping 冲突、权限错误、字段格式错误。不可重试失败如果一直自动重试,只会制造消息堆积。
消费端不要吞异常
错误写法:
try {
productSearchWriter.upsert(doc);
} catch (Exception ex) {
log.error("sync es failed", ex);
}这段代码的问题是:日志打了,但消息被当成消费成功,MQ 不会再投递,ES 就可能永久停留在旧数据。
推荐写法:
try {
productSearchWriter.upsert(doc);
syncLogRepository.markSuccess(event.getEventId());
} catch (Exception ex) {
syncLogRepository.markFailed(event.getEventId(), ex.getMessage());
throw ex;
}抛出异常后,MQ 才能触发重试或进入死信。同步日志也能让排查人员看到失败事件、业务 ID、版本号和错误原因。
重试要退避,不能无限猛冲
ES 更新失败后不能立刻无限重试。因为 ES 已经抖动时,猛烈重试会继续放大压力。
常见退避策略:
| 第几次失败 | 下次重试时间 |
|---|---|
| 第 1 次 | 10 秒后 |
| 第 2 次 | 30 秒后 |
| 第 3 次 | 1 分钟后 |
| 第 4 次 | 5 分钟后 |
| 第 5 次 | 15 分钟后 |
| 超过阈值 | 进入死信或人工处理 |
伪代码:
public Date nextRetryTime(int retryCount) {
int[] seconds = new int[] {10, 30, 60, 300, 900};
int index = Math.min(retryCount, seconds.length - 1);
return new Date(System.currentTimeMillis() + seconds[index] * 1000L);
}退避的目的不是“慢”,而是给 ES 恢复时间,同时避免把故障从搜索系统扩散到 MQ、消费者和数据库。
Bulk 更新必须检查每一条
Bulk 请求返回成功,不代表每条文档都成功。必须看 errors 和每个 item 的状态。
错误思路:
HTTP 200 = 全部成功正确思路:
HTTP 200 只代表 Bulk 请求被 ES 接收,里面每条 index/update/delete 都可能单独成功或失败。处理流程:
flowchart TD
A["发送 Bulk"] --> B{"Bulk 是否整体失败"}
B -- "是" --> C["整批按退避重试"]
B -- "否" --> D["检查每个 item"]
D --> E{"是否有失败 item"}
E -- "否" --> F["整批成功"]
E -- "是" --> G["记录失败 bizId 和原因"]
G --> H["只重试失败 item"]Java 伪代码:
BulkResult result = productSearchWriter.bulkUpsert(docs);
if (result.hasErrors()) {
for (BulkItemResult item : result.getItems()) {
if (!item.isSuccess()) {
syncLogRepository.markFailed(
item.getEventId(),
item.getBizId(),
item.getErrorMessage()
);
retryRepository.save(item.getEventId(), item.getBizId(), nextRetryTime(0));
}
}
}真实项目中,要把失败项的 index、id、status、error.type、error.reason 记录下来。否则后面只知道“同步失败”,不知道失败在哪一条。
死信不是结束,是待处理区
消息进入死信队列,不代表可以不管。死信的含义是:自动重试已经解决不了,需要告警、定位和人工或任务重放。
死信处理要包含:
| 动作 | 目的 |
|---|---|
| 记录 bizId、eventId、version | 后续能精准重放 |
| 记录失败原因 | 区分 ES 故障、Mapping 冲突、数据脏值 |
| 告警 | 防止失败长期没人知道 |
| 支持单条重放 | 修复少量问题 |
| 支持按时间范围重放 | 修复一批问题 |
| 支持跳过或标记无效 | 处理确实不应再同步的数据 |
死信重放时仍然要查 MySQL 最新快照,不要直接拿死信里的旧文档写 ES。否则死信越晚处理,越可能把旧状态写回去。
更新失败后的数据状态怎么理解
ES 更新失败后,数据状态可能是下面几种:
| 状态 | 说明 | 用户可能看到什么 |
|---|---|---|
| MySQL 新,ES 旧 | 最常见 | 搜索列表显示旧价格、旧状态 |
| MySQL 有,ES 没有 | 新增同步失败 | 搜索不到 |
| MySQL 删除,ES 还在 | 删除同步失败 | 下架或删除数据仍能搜到 |
| ES 写了一半字段 | 局部更新失败或脚本错误 | 字段不一致 |
| Bulk 部分成功 | 一批数据中部分新、部分旧 | 同一页搜索结果新旧混杂 |
所以搜索结果不能作为交易事实。下单价格、库存、权限、最终状态必须回业务服务或 MySQL 校验。
补偿任务兜底
即使 MQ 重试和死信都做了,也还要有补偿。因为可能存在未知失败:代码 bug、人工改库、历史数据、消费者逻辑漏处理。
补偿任务可以这样做:
flowchart TD
A["扫描 MySQL 最近变更数据"] --> B["按 ID 查询 ES 文档"]
B --> C{"version 是否一致"}
C -- "一致" --> D["跳过"]
C -- "不一致" --> E["查 MySQL 最新快照"]
E --> F["重新 upsert 或 delete ES"]
F --> G["记录补偿结果"]补偿任务重点比较:
| 字段 | 为什么比较 |
|---|---|
version | 快速判断 ES 是否落后 |
updatedAt | 判断同步延迟 |
status | 下架、删除、禁用最容易出事故 |
| 核心展示字段 | 商品名、价格、资产名、标签等 |
searchVisible | 判断是否应该被搜索到 |
项目里推荐的处理闭环
一套完整处理闭环应该是:
- MySQL 先提交,ES 不参与主事务。
- 通过本地消息表、事务消息或 CDC 产生可靠事件。
- 消费端查 MySQL 最新快照。
- 写 ES 成功后记录成功。
- 写 ES 失败后记录失败并抛异常。
- MQ 按退避策略重试。
- 超过重试次数进入死信。
- 死信告警并支持人工重放。
- 定时补偿扫描 MySQL 和 ES 差异。
- 关键业务读取 MySQL 或业务服务,不依赖 ES 搜索结果。
线上排查流程
商品改价后 ES 还是旧价格
flowchart TD
A["搜索页显示旧价格"] --> B["查 MySQL 当前价格"]
B --> C["按 ID 查 ES 文档"]
C --> D{"ES 是否旧值"}
D -- "否" --> E["检查搜索 DSL 或缓存"]
D -- "是" --> F["查同步日志和 MQ 延迟"]
F --> G{"是否有失败或堆积"}
G -- "是" --> H["重试、死信重放、扩容消费者"]
G -- "否" --> I["查乱序、版本、补偿任务"]排查清单:
| 检查项 | 命令或方法 |
|---|---|
| MySQL 当前值 | 按主键查业务表 |
| ES 当前文档 | GET /index/_doc/id |
| 搜索 DSL | 打印最终 DSL,确认过滤和排序 |
| 同步消息 | 查 MQ 是否堆积、是否死信 |
| 同步日志 | 查 eventId、bizId、失败原因 |
| 版本字段 | 比较 MySQL version 和 ES version |
| refresh | 判断写入后是否已过 refresh 窗口 |
| 补偿任务 | 查最近补偿是否成功 |
下架商品仍然能搜到
重点查:
- MySQL 状态是否已经下架。
- ES 文档
searchVisible或status是否同步。 - 搜索 DSL 是否过滤
searchVisible=true。 - 删除或下架消息是否失败。
- 是否有旧更新消息把文档重新写回。
主库有数据 ES 搜不到
重点查:
- 是否同步到 ES。
- ES 文档是否存在。
- refresh 是否完成。
- Mapping 是否正确。
term和match是否用错。- 分词结果是否符合预期。
- filter 是否误伤。
- 权限或租户条件是否错误。
监控指标
| 指标 | 说明 |
|---|---|
| 同步延迟 | MySQL updated_at 到 ES updatedAt 的差值 |
| MQ 堆积 | 变更事件积压数量 |
| 消费失败数 | 写 ES 失败、组装失败、网络异常 |
| 死信数量 | 长期无法成功的消息 |
| 补偿修复数 | 定时任务发现并修复的差异 |
| 对账差异率 | 抽样或全量差异比例 |
| ES 写入耗时 | upsert/delete 的耗时 |
| ES rejected | 写入线程池拒绝数 |
| refresh 延迟 | 搜索可见性延迟 |
没有监控的一致性方案,本质上是靠运气。
常见错误
| 错误 | 后果 | 正确做法 |
|---|---|---|
| 写 MySQL 后同步写 ES | ES 抖动拖慢主流程,失败难补 | MQ/CDC 异步 + 补偿 |
| 消息里放完整旧数据 | 乱序时旧数据覆盖新状态 | 消费端查主库最新快照 |
| 消费不做幂等 | 重复消息导致重复或异常 | eventId、文档 ID、唯一约束 |
| 删除只改 MySQL | ES 残留旧文档 | 删除或标记不可搜也要同步 |
| 搜索结果用于下单价格 | 价格可能短暂旧 | 下单回商品服务或 MySQL 校验 |
| 没有补偿任务 | 一次失败可能长期不一致 | 重试、死信、定时对账 |
| 重建索引不追增量 | 切换后新索引缺数据 | 全量 + 增量追平 + 校验 |
| 只校验文档数量 | 字段、分词、状态仍可能错 | 数量、抽样、关键查询一起校验 |
商业场景:医疗数据采集与资产平台
资产平台里,MySQL 可以保存资产、表、字段、任务、机构、标签等事实数据;ES 保存面向检索的资产视图。
典型文档:
{
"assetId": 10001,
"assetName": "门诊检验结果表",
"orgId": 2001,
"orgName": "第一人民医院",
"tableName": "ods_lab_result",
"fieldNames": ["patient_id", "test_item", "result_value"],
"tags": ["检验", "门诊", "患者主索引"],
"status": "ONLINE",
"version": 88,
"updatedAt": "2026-07-04T10:00:00"
}一致性设计:
- 资产编辑、上下线、标签变更先提交 MySQL。
- 本地消息表或 MQ 产生
ASSET_CHANGED事件。 - 同步服务按资产 ID 查询最新资产快照和字段列表。
- 组装 ES 文档并 upsert。
- 资产下线后设置
searchVisible=false或删除文档。 - 定时补偿扫描最近变更资产,对比 ES
version。 - 搜索结果点击详情时,以 MySQL 详情页为准。
关联知识点
| 知识点 | 继续学习什么 |
|---|---|
| 同步与重建索引 | MQ、CDC、重建索引、别名切换 |
| 底层原理 | refresh、segment、更新和近实时 |
| 性能优化与排查 | 搜不到、搜不准、同步慢、写入慢 |
| 可靠消息最终一致 | 本地消息表、Outbox、事务消息 |
| CDC与Outbox数据同步 | MySQL 到 ES、Redis、MQ 的位点、重复、乱序、重放和对账 |
| 幂等设计 | 重复消息和重复请求怎么处理 |
| 数据一致性 | 强一致、最终一致、补偿和对账 |
| 消息堆积与背压 | 同步消息堆积后的排查 |
本章小结
MySQL 与 ES 一致性的核心不是追求“同时成功”,而是让 MySQL 先成为可靠事实源,再用可靠事件把变化同步到 ES。同步链路必须处理失败、重复、乱序、删除、refresh 延迟、重建期间增量和长期差异。商业项目的成熟做法是:本地消息表或 CDC 保证不丢变更,消费端查主库最新快照,ES 写入幂等且有版本控制,失败进入重试和死信,定时补偿和对账兜底,关键业务读 MySQL,搜索列表读 ES。
