Skip to content

MySQL 与 ES 数据一致性

商业系统里,MySQL 通常是事实数据源,Elasticsearch 是搜索视图。事实数据源负责事务、约束、状态机和最终正确性;搜索视图负责全文检索、聚合筛选、相关性排序和高效查询。

一句话理解:

MySQL 和 ES 很难做强一致,商业项目通常做最终一致:主库先正确提交,再通过可靠同步、幂等消费、乱序控制、补偿任务、对账校验,把 ES 修正到和主库一致。

学习目标

目标需要掌握什么
知道定位MySQL 是事实源,ES 是搜索副本
知道风险双写失败、消息丢失、重复消费、乱序、删除漏同步、refresh 延迟
知道方案同步双写、MQ、Binlog CDC、本地消息表、事务消息、定时补偿
知道原理为什么最终一致,为什么消费端要查主库最新快照
会写 Demo表设计、事件消息、消费者 upsert/delete、补偿任务
会排查搜不到、搜到旧数据、下架仍可搜、同步延迟、消息堆积

如果你想把这一页的同步、一致性、失败补偿和重建流程放到完整商业搜索链路里练习,继续做:Elasticsearch 商业场景训练营

为什么不能要求强一致

MySQL 和 ES 是两个不同系统。一次业务写入如果同时写 MySQL 和 ES,就会遇到典型分布式一致性问题:

mermaid
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 失败可以自动重试
幂等同一事件重复消费不会写错
防乱序旧事件不能覆盖新状态
可补偿已知失败和未知漏同步都能修复
可观测能看到延迟、失败、积压、差异

不要承诺:

  1. MySQL 提交后 ES 立即可搜索。
  2. ES 搜索结果永远等于 MySQL 当前状态。
  3. 消息只投递一次。
  4. 消费顺序永远和业务写入顺序一致。

常见同步方案对比

方案流程优点风险推荐程度
同步双写写 MySQL 后立即写 ES简单主流程受 ES 影响,失败难补
MQ 异步写库后发消息,消费者写 ES解耦、可重试有延迟,要防丢和乱序
本地消息表业务数据和消息表同事务提交解决写库成功但消息没发多一张表和投递任务
RocketMQ 事务消息半消息 + 本地事务 + 提交消息减少消息与事务不一致依赖 MQ 能力和回查中高
Binlog CDC监听 MySQL binlog 同步业务侵入小跨表组装复杂,链路复杂
定时补偿扫主库对比 ES 并修复兜底可靠延迟高必须有
全量重建从主库重建新索引修复结构和历史差异成本高结构变更时用

成熟方案通常是:

mermaid
flowchart TD
    A["MySQL 本地事务"] --> B["变更事件"]
    B --> C["MQ 或 CDC"]
    C --> D["同步服务写 ES"]
    D --> E["失败重试和死信"]
    F["定时补偿"] --> D
    G["对账校验"] --> F

一致性链路要分别保证什么

MySQL 到 ES 的同步不是“把消息发出去”这么简单,而是一条端到端链路。任何一段没有兜底,最后都会表现成“搜索页看到旧数据”。

mermaid
flowchart TD
    A["业务写入"] --> B["事件产生"]
    B --> C["事件投递"]
    C --> D["事件消费"]
    D --> E["写入 ES"]
    E --> F["搜索可见"]
    F --> G["补偿对账"]

每一段要解决的问题不同:

链路要保证什么常见手段如果不处理会怎样
业务写入MySQL 数据先正确提交本地事务、唯一约束、状态机校验主数据本身就是错的,ES 同步再准也没意义
事件产生主库变更一定有同步事件本地消息表、事务消息、CDCMySQL 成功但没有消息,ES 永远旧
事件投递事件能送到同步服务MQ 重试、确认机制、死信、投递日志消息丢失或长期堆积
事件消费重复和乱序不会写错幂等表、业务主键、版本号、按 key 顺序重复消费报错,旧事件覆盖新状态
写入 ES写失败可以重试Bulk、失败明细、退避重试、死信重放单条失败混在批次里没人发现
搜索可见写入后理解 refresh 延迟refresh 窗口、refresh=wait_for、读 MySQL 兜底误把近实时延迟当成同步失败
补偿对账未知差异能被发现和修复定时扫描、抽样对账、版本对比、人工重放小概率错误沉淀成长期脏数据

所以“数据一致性”不是一个 API,而是一套工程机制:可靠事件、幂等消费、乱序控制、失败重试、补偿对账、读路径兜底。

为什么同步双写不推荐

同步双写看起来最简单:

java
public void updateProduct(Product product) {
    productMapper.update(product);
    elasticsearchClient.update(product);
}

但它有很多失败窗口:

mermaid
flowchart TD
    A["开始修改商品"] --> B["MySQL 更新成功"]
    B --> C{"写 ES 是否成功"}
    C -- "成功" --> D["暂时一致"]
    C -- "失败" --> E["MySQL 新,ES 旧"]
    C -- "超时" --> F["不知道 ES 是否成功"]

同步双写的问题:

  1. ES 抖动会拖慢主业务接口。
  2. MySQL 成功、ES 失败后需要补偿。
  3. ES 超时时不知道到底写没写成功。
  4. 如果业务事务后来回滚,ES 可能已经写入错误数据。
  5. 高并发更新时容易发生旧请求覆盖新请求。

如果系统很小、后台低频、允许人工修复,可以临时用同步双写;核心商业链路不建议。

失败窗口怎么分析

设计一致性方案时,不要只看“正常流程”,要按失败窗口推演。

失败点可能结果解决办法
MySQL 事务提交前服务宕机业务数据没提交,事件也不应产生本地事务保证一起回滚
MySQL 提交成功但还没发 MQ 就宕机主库新,ES 旧本地消息表保存待发送事件,后台继续投递
MQ 已收到但业务事务回滚ES 可能同步了不存在或错误状态先提交事务再投递,或使用事务消息回查
MQ 重复投递同一事件被消费多次event_id 唯一约束和 ES 业务 _id 幂等
消息乱序旧价格覆盖新价格,下架后又可搜索消费端查主库最新快照,或使用 version 防旧覆盖
ES Bulk 部分失败一批里有的成功有的失败逐条检查 Bulk item,失败明细入库重试
消费者一直失败消息堆积,ES 长期旧重试退避、死信队列、告警、人工重放
全量重建期间有增量新索引切换后缺最新数据记录起点,增量追平,校验后切别名

能把这些失败窗口讲清楚,才算真正理解 MySQL 和 ES 的一致性。

RocketMQ 事务消息怎么用在同步里

RocketMQ 事务消息可以降低“本地事务和消息发送不一致”的风险。它不是让 MySQL 和 ES 强一致,而是尽量保证“业务成功后消息最终可见”。

mermaid
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

本地消息表解决的是“业务提交成功,但消息没发出去”的问题。

mermaid
flowchart TD
    A["业务服务"] --> B["开启 MySQL 事务"]
    B --> C["更新业务表"]
    C --> D["插入 outbox 事件表"]
    D --> E["提交事务"]
    E --> F["投递任务发送 MQ"]
    F --> G["同步服务消费"]
    G --> H["写 ES"]

业务数据和事件表在同一个 MySQL 本地事务里提交,要么都成功,要么都失败。

表设计 Demo

商品表:

sql
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;

事件表:

sql
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下次重试时间

写业务和事件

java
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 包住业务更新和事件插入。

投递任务

java
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));
        }
    }
}

投递任务要支持:

  1. 批量拉取。
  2. 失败重试。
  3. 最大重试次数。
  4. 告警。
  5. 人工重投。

推荐架构二:Binlog CDC

CDC Change Data Capture 的思路是监听 MySQL binlog,把数据库变更转换为同步事件。常见工具有 Canal、Debezium、Flink CDC 等。

mermaid
flowchart TD
    A["业务服务写 MySQL"] --> B["MySQL binlog"]
    B --> C["CDC 组件"]
    C --> D["变更事件"]
    D --> E["同步服务"]
    E --> F["组装 ES 文档"]
    F --> G["写 ES"]

CDC 的优点:

  1. 对业务代码侵入小。
  2. 不需要每个业务方法手动发消息。
  3. 能捕获直接 SQL 修改带来的变更。
  4. 适合数据平台、搜索平台统一同步。

CDC 的难点:

难点说明
跨表组装ES 文档常来自多张表,单表 binlog 不一定够
变更顺序要按主键或业务维度处理乱序
删除事件要正确处理 delete 和逻辑删除
全量加增量第一次初始化要处理全量期间的增量
位点保存CDC 消费位点要可靠保存
字段语义binlog 只有字段变化,不一定知道业务意图

CDC 适合“统一数据同步平台”,业务消息适合“业务语义清楚的事件驱动”。很多商业项目会二者结合。

消费端为什么要查 MySQL 最新快照

同步消息有两种设计:

方式内容风险
消息放完整文档商品名、价格、库存、标签全放消息里消息大,字段变化难维护,乱序风险高
消息只放 ID 和版本消费时查 MySQL 最新状态多一次查库,但更稳

推荐消费端按 ID 查询 MySQL 最新快照:

mermaid
flowchart TD
    A["消费商品变更事件"] --> B["按商品 ID 查 MySQL 最新快照"]
    B --> C{"当前是否可搜索"}
    C -- "是" --> D["组装最新 ES 文档"]
    D --> E["upsert 到 ES"]
    C -- "否" --> F["删除或标记不可搜"]

这样可以减少乱序问题。比如商品先改价,再下架。如果“改价消息”晚于“下架消息”被消费,消费端查 MySQL 最新状态发现商品已经下架,就不会把旧价格文档重新写回可搜索状态。

幂等设计

MQ 和 CDC 都可能重复投递。消费端必须做到重复执行结果一致。

文档 ID 幂等

ES 文档 _id 使用业务主键:

json
PUT /product_search/_doc/10001
{
  "id": 10001,
  "productName": "无线蓝牙耳机",
  "version": 101
}

同一个商品重复写入同一个 _id,不会生成多份文档。

消费日志幂等

同步日志表:

sql
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 同步最容易被问到的问题。

典型场景:

text
10:00:01 商品改价 version=101
10:00:02 商品下架 version=102

如果 version=102 先写 ES,version=101 后写 ES,就可能旧数据覆盖新数据。

解决方式:

方式思路适合场景
消费时查主库最新快照永远以 MySQL 当前状态组装文档最常用
业务版本号旧版本不能覆盖新版本高频更新
按业务 key 顺序消费同一个商品进入同一分区Kafka、RocketMQ 顺序场景
ES 外部版本用版本控制写入对版本管理清晰的场景

版本号脚本 Demo

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.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 顺序消费”。

删除和下架一致性

删除比更新更容易出事故。

常见错误:

  1. MySQL 已删除,ES 文档没删,搜索还能搜到。
  2. 商品下架只是改状态,但 ES 查询没有过滤状态。
  3. 物理删除消息乱序后,旧更新又把文档写回来。
  4. 全量重建时把已删除数据重新导入。

推荐策略:

业务动作ES 处理
商品下架设置 searchVisible=false 或删除文档
商品软删除设置删除状态,并搜索时过滤
商品物理删除删除 ES 文档
订单状态变化更新状态字段,但订单事实仍查 MySQL

逻辑不可搜:

json
POST /product_search/_update/10001
{
  "doc": {
    "searchVisible": false,
    "status": "OFF_SALE",
    "version": 103,
    "updatedAt": "2026-07-04T10:00:03"
  }
}

搜索 DSL 必须过滤:

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

物理删除:

json
DELETE /product_search/_doc/10001

删除也要幂等。重复删除同一个文档应该被视为成功或可忽略。

Refresh 延迟不是同步失败

ES 写入成功后,不代表立刻能被搜索到。只有 refresh 生成可搜索 segment 后,搜索请求才能看到新文档。

mermaid
flowchart TD
    A["写入 ES 成功"] --> B["进入内存 buffer 和 translog"]
    B --> C["等待 refresh"]
    C --> D["生成可搜索 segment"]
    D --> E["搜索可见"]

常见误判:

现象可能原因
按 ID GET 能查到,搜索查不到refresh 未完成或查询条件不匹配
写入后 1 秒内搜不到近实时正常窗口
长时间搜不到同步失败、Mapping、分词、过滤条件问题

如果某个后台操作要求“提交后立刻可见”,可以:

  1. 操作完成页直接读 MySQL。
  2. 给用户展示“同步中”状态。
  3. 对少量关键写入使用 refresh=wait_for
  4. 不要全局把 refresh_interval 调得极短,否则写入和 merge 压力会上升。

读路径怎么设计

搜索列表读 ES,详情和交易校验读 MySQL。

mermaid
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、人工改库、历史脏数据

按更新时间补偿

mermaid
flowchart TD
    A["定时扫描 MySQL 最近变更"] --> B["查询对应 ES 文档"]
    B --> C{"关键字段是否一致"}
    C -- "一致" --> D["跳过"]
    C -- "不一致" --> E["重新写入 ES"]
    E --> F["记录补偿日志"]

SQL:

sql
select id, version, updated_at
from product
where updated_at >= ?
  and updated_at < ?
order by updated_at asc
limit 1000;

补偿要注意:

  1. 分批扫描,避免压垮主库。
  2. 按时间窗口推进,保存游标。
  3. 对失败记录重试。
  4. 对长期失败告警。
  5. 不要和全量重建互相冲突。

抽样对账

对账不一定每次全量,可以抽样或按业务范围:

对账方式适合场景
最近 1 小时变更全量对比发现同步延迟和失败
每天抽样 1 万个 ID发现长期差异
按租户或机构分批对账多租户数据隔离
按状态对账下架、删除、异常状态重点查
全量对账大版本发布或重大修复后

全量重建时的一致性

Mapping、分词器、字段类型改了,通常要重建索引。重建期间也会有业务增量变更,如果处理不好,新索引切过去就是旧数据。

推荐流程:

mermaid
flowchart TD
    A["创建新索引 v2"] --> B["记录全量开始时间 T1"]
    B --> C["从 MySQL 全量导入 v2"]
    C --> D["回放 T1 之后的增量"]
    D --> E["校验数量和抽样字段"]
    E --> F{"校验通过"}
    F -- "否" --> G["修复并重新追平"]
    F -- "是" --> H["原子切换别名"]

关键点:

  1. 业务访问别名,不写死版本索引。
  2. 全量导入期间的增量不能丢。
  3. 切换前必须校验。
  4. 切换后保留旧索引,方便回滚。
  5. 如果使用 CDC,要保存 binlog 位点;如果使用 MQ,要保存事件时间或版本。

Java 消费者 Demo

下面示例使用 JDK 8 风格,避免依赖高版本语法。

java
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 体现几个原则:

  1. eventId 做消费幂等。
  2. 消费时查 MySQL 最新快照。
  3. 当前不可搜索就删除或标记 ES 文档。
  4. 失败不吞异常,让 MQ 重试或进入死信。
  5. 同步日志落库,便于排查。

ES 更新失败怎么办

ES 更新失败时,第一原则是:不能影响 MySQL 主事务的正确提交,也不能把失败静默吞掉。因为 MySQL 是事实源,ES 是搜索视图。ES 短暂失败可以接受,长期失败必须被发现、重试、补偿和告警。

整体处理流程:

mermaid
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 冲突、权限错误、字段格式错误。不可重试失败如果一直自动重试,只会制造消息堆积。

消费端不要吞异常

错误写法:

java
try {
    productSearchWriter.upsert(doc);
} catch (Exception ex) {
    log.error("sync es failed", ex);
}

这段代码的问题是:日志打了,但消息被当成消费成功,MQ 不会再投递,ES 就可能永久停留在旧数据。

推荐写法:

java
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 分钟后
超过阈值进入死信或人工处理

伪代码:

java
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 的状态。

错误思路:

text
HTTP 200 = 全部成功

正确思路:

text
HTTP 200 只代表 Bulk 请求被 ES 接收,里面每条 index/update/delete 都可能单独成功或失败。

处理流程:

mermaid
flowchart TD
    A["发送 Bulk"] --> B{"Bulk 是否整体失败"}
    B -- "是" --> C["整批按退避重试"]
    B -- "否" --> D["检查每个 item"]
    D --> E{"是否有失败 item"}
    E -- "否" --> F["整批成功"]
    E -- "是" --> G["记录失败 bizId 和原因"]
    G --> H["只重试失败 item"]

Java 伪代码:

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));
        }
    }
}

真实项目中,要把失败项的 indexidstatuserror.typeerror.reason 记录下来。否则后面只知道“同步失败”,不知道失败在哪一条。

死信不是结束,是待处理区

消息进入死信队列,不代表可以不管。死信的含义是:自动重试已经解决不了,需要告警、定位和人工或任务重放。

死信处理要包含:

动作目的
记录 bizId、eventId、version后续能精准重放
记录失败原因区分 ES 故障、Mapping 冲突、数据脏值
告警防止失败长期没人知道
支持单条重放修复少量问题
支持按时间范围重放修复一批问题
支持跳过或标记无效处理确实不应再同步的数据

死信重放时仍然要查 MySQL 最新快照,不要直接拿死信里的旧文档写 ES。否则死信越晚处理,越可能把旧状态写回去。

更新失败后的数据状态怎么理解

ES 更新失败后,数据状态可能是下面几种:

状态说明用户可能看到什么
MySQL 新,ES 旧最常见搜索列表显示旧价格、旧状态
MySQL 有,ES 没有新增同步失败搜索不到
MySQL 删除,ES 还在删除同步失败下架或删除数据仍能搜到
ES 写了一半字段局部更新失败或脚本错误字段不一致
Bulk 部分成功一批数据中部分新、部分旧同一页搜索结果新旧混杂

所以搜索结果不能作为交易事实。下单价格、库存、权限、最终状态必须回业务服务或 MySQL 校验。

补偿任务兜底

即使 MQ 重试和死信都做了,也还要有补偿。因为可能存在未知失败:代码 bug、人工改库、历史数据、消费者逻辑漏处理。

补偿任务可以这样做:

mermaid
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判断是否应该被搜索到

项目里推荐的处理闭环

一套完整处理闭环应该是:

  1. MySQL 先提交,ES 不参与主事务。
  2. 通过本地消息表、事务消息或 CDC 产生可靠事件。
  3. 消费端查 MySQL 最新快照。
  4. 写 ES 成功后记录成功。
  5. 写 ES 失败后记录失败并抛异常。
  6. MQ 按退避策略重试。
  7. 超过重试次数进入死信。
  8. 死信告警并支持人工重放。
  9. 定时补偿扫描 MySQL 和 ES 差异。
  10. 关键业务读取 MySQL 或业务服务,不依赖 ES 搜索结果。

线上排查流程

商品改价后 ES 还是旧价格

mermaid
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 窗口
补偿任务查最近补偿是否成功

下架商品仍然能搜到

重点查:

  1. MySQL 状态是否已经下架。
  2. ES 文档 searchVisiblestatus 是否同步。
  3. 搜索 DSL 是否过滤 searchVisible=true
  4. 删除或下架消息是否失败。
  5. 是否有旧更新消息把文档重新写回。

主库有数据 ES 搜不到

重点查:

  1. 是否同步到 ES。
  2. ES 文档是否存在。
  3. refresh 是否完成。
  4. Mapping 是否正确。
  5. termmatch 是否用错。
  6. 分词结果是否符合预期。
  7. filter 是否误伤。
  8. 权限或租户条件是否错误。

监控指标

指标说明
同步延迟MySQL updated_at 到 ES updatedAt 的差值
MQ 堆积变更事件积压数量
消费失败数写 ES 失败、组装失败、网络异常
死信数量长期无法成功的消息
补偿修复数定时任务发现并修复的差异
对账差异率抽样或全量差异比例
ES 写入耗时upsert/delete 的耗时
ES rejected写入线程池拒绝数
refresh 延迟搜索可见性延迟

没有监控的一致性方案,本质上是靠运气。

常见错误

错误后果正确做法
写 MySQL 后同步写 ESES 抖动拖慢主流程,失败难补MQ/CDC 异步 + 补偿
消息里放完整旧数据乱序时旧数据覆盖新状态消费端查主库最新快照
消费不做幂等重复消息导致重复或异常eventId、文档 ID、唯一约束
删除只改 MySQLES 残留旧文档删除或标记不可搜也要同步
搜索结果用于下单价格价格可能短暂旧下单回商品服务或 MySQL 校验
没有补偿任务一次失败可能长期不一致重试、死信、定时对账
重建索引不追增量切换后新索引缺数据全量 + 增量追平 + 校验
只校验文档数量字段、分词、状态仍可能错数量、抽样、关键查询一起校验

商业场景:医疗数据采集与资产平台

资产平台里,MySQL 可以保存资产、表、字段、任务、机构、标签等事实数据;ES 保存面向检索的资产视图。

典型文档:

json
{
  "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"
}

一致性设计:

  1. 资产编辑、上下线、标签变更先提交 MySQL。
  2. 本地消息表或 MQ 产生 ASSET_CHANGED 事件。
  3. 同步服务按资产 ID 查询最新资产快照和字段列表。
  4. 组装 ES 文档并 upsert。
  5. 资产下线后设置 searchVisible=false 或删除文档。
  6. 定时补偿扫描最近变更资产,对比 ES version
  7. 搜索结果点击详情时,以 MySQL 详情页为准。

关联知识点

知识点继续学习什么
同步与重建索引MQ、CDC、重建索引、别名切换
底层原理refresh、segment、更新和近实时
性能优化与排查搜不到、搜不准、同步慢、写入慢
可靠消息最终一致本地消息表、Outbox、事务消息
CDC与Outbox数据同步MySQL 到 ES、Redis、MQ 的位点、重复、乱序、重放和对账
幂等设计重复消息和重复请求怎么处理
数据一致性强一致、最终一致、补偿和对账
消息堆积与背压同步消息堆积后的排查

本章小结

MySQL 与 ES 一致性的核心不是追求“同时成功”,而是让 MySQL 先成为可靠事实源,再用可靠事件把变化同步到 ES。同步链路必须处理失败、重复、乱序、删除、refresh 延迟、重建期间增量和长期差异。商业项目的成熟做法是:本地消息表或 CDC 保证不丢变更,消费端查主库最新快照,ES 写入幂等且有版本控制,失败进入重试和死信,定时补偿和对账兜底,关键业务读 MySQL,搜索列表读 ES。