在公司里做后端开发,但凡业务里同时用到了MySQL和Elasticsearch(以下简称ES),十有八九会被同一个问题缠住:两边数据对不上怎么办?这是个经典问题,网上讨论也很多,但大部分文章要么停在“用消息队列”这种口号上,要么只讲Canal怎么搭,缺少一篇把“为什么不一致”“有哪些路可以走”“每条路的坑在哪”捋清楚的总结。我这些年经手过好几个电商、内容、检索类的项目,在MySQL和ES同步上踩过不少坑,也沉淀了一些自己的判断。这篇就把我对这个问题的完整认知写出来,包括方案对比、落地细节、兜底策略和选型建议,希望对正在被这个问题折磨的同行有帮助。
1. 先搞清楚“一致性”到底指什么
1.1 架构真相:MySQL和ES本来就不是一个世界的
在讨论怎么保证一致性之前,得先把问题定义清楚。MySQL是关系型数据库,支持事务、强一致、行级锁,数据是业务的“唯一事实来源”;ES是分布式搜索引擎,基于Lucene,擅长全文检索和聚合分析,它的索引本质上是“倒排索引”,数据是“副本”性质的。
这两套系统为什么需要同步?因为MySQL的查询能力跟不上业务需求——模糊搜索、分词、复杂聚合、地理位置检索,这些是ES的主场。但ES不擅长事务,不支持跨文档更新,也没有像MySQL那样成熟的MVCC机制。于是几乎所有团队都选择“MySQL写主库,ES建索引”的架构。
可问题来了:MySQL的写入是实时的、事务性的,而ES的写入是异步的、近实时的(默认refresh间隔1秒)。这两者之间存在天然的时间窗口,任何同步方案都不可能做到100%零延迟,所以一致性问题的本质是:在MySQL和ES之间存在天然的异步窗口,我们要做的是把这个窗口控制到最小,并且保证窗口之外的数据绝对一致。
1.2 不一致的现象、根因与代价
我把项目中实际遇到的不一致现象整理了一下,大致分三类:
第一类是延迟性不一致。MySQL里改了一条数据,ES还没收到,搜索出来的是旧值。这类现象最普遍,用户感知也最强——刚改完商品标题,搜索结果还是旧的。
第二类是丢失性不一致。MySQL更新成功,但同步链路断了(消费者挂了、MQ丢了、代码异常),ES永远停留在旧状态。这类问题最致命,必须靠兜底机制发现和修复。
第三类是错乱性不一致。数据的更新顺序颠倒,比如先发了更新A、再发了更新B,但消费者先处理了B后处理了A,导致ES里存的不是最新值。这类问题隐蔽,排查起来很费劲。
根本原因通常逃不过这几个:业务代码里“双写”逻辑不完整,只写了MySQL忘了写ES;异步任务里没有重试机制,一旦失败就静默吞掉;MQ乱序或重复投递没有处理;全量同步和增量同步的冲突没控制好;线上变更(比如ES索引重建、分片迁移)期间切流处理不当。
代价有多大呢?我见到的真实案例,一个电商平台因为商品价格在ES里没更新,导致用户搜索出来的价格和详情页不一致,被投诉到客服;一个内容社区因为文章状态没同步,让用户搜到了已删除的文章。这些都是线上事故级别的,所以一致性不是“锦上添花”的优化项,而是必须正面解决的工程问题——这也是我为什么把方案分成“同步链路方案”和“兜底机制”两层,前者负责尽量实时,后者负责最终兜底。
2. 常见同步方案盘点:从“能用”到“好用”,坑在哪里
2.1 方案速览与选型地图
业界主流的MySQL同步到ES方案,大致可以分成四类:业务双写、定时任务扫表、Binlog订阅(Canal/Debezium)、基于Binlog + MQ的异步框架。每一类都有自己的适用场景,没有银弹。
为了让你一目了然,我先给一张对比表,后面再逐层展开。
| 方案 | 实时性 | 数据可靠性 | 侵入性 | 实现成本 | 适用场景 |
|---|---|---|---|---|---|
| 业务双写 | 高 | 低(易漏写) | 高 | 低 | 小项目、演示项目 |
| 定时任务扫表 | 低(分钟级) | 中 | 低 | 低 | 数据量小、容忍延迟 |
| Canal订阅Binlog | 高 | 高 | 低 | 中 | 增量同步、与业务解耦 |
| Binlog + MQ | 高 | 高 | 低 | 中高 | 数据量大、要求最终一致 |
表格列完了,逐个说重点。业务双写就是业务代码里在写完MySQL之后,再手动调ES的写接口。这是最直观的思路,也是新手最容易踩坑的方案。你以为写完了就完了?事务只包裹了MySQL,ES写失败了怎么办?库存扣减成功了,ES里的库存没变,用户搜出来还有货,下单却失败。双写方案最大的问题是业务代码里遍布ES逻辑,极度耦合,而且MySQL事务和ES异步写天然无法原子,失败场景没法优雅处理。
定时任务扫表一般是写个Job每隔几分钟把变更的数据捞出来更新ES。实现很简单,但是延迟分钟级是常态,只适合对实时性要求不高的数据。但这套方案有个关键优势——可以作为兜底手段,后面专门讲。
Canal订阅Binlog是目前最主流的增量同步方案。Canal伪装成MySQL的从库,订阅Binlog事件,解析成结构化数据,再推送给下游。它好在完全不用碰业务代码,对业务无侵入,还可以精准感知到行的变更(INSERT/UPDATE/DELETE)。但是只用Canal直接推给ES的话,面临两个问题:一是Canal消费速度跟ES写入速度之间可能不匹配,Canal是连续读Binlog的,一旦ES写入变慢或故障,Canal这边会积压甚至丢数据;二是没有消息持久化,消费者重启后没法从断点续传。所以单用Canal搭同步链路,稳定性上不太够。
Binlog + MQ就是在Canal和ES之间加一层消息队列,比如RocketMQ或Kafka。这样Canal负责把Binlog转成消息,MQ负责削峰填谷、持久化、回溯消费,消费者按需把数据写进ES。这是我认为生产环境最值得推荐的方案,后面我会重点展开。
2.2 为什么推荐“Binlog + MQ”而不是“双写 + MQ”
有人可能会说,我不是用MQ吗?我在业务代码里同步发个MQ,消费者消费了再写ES,不也一样解耦吗?这个思路看起来对,但仔细推敲会发现两个致命隐患。
第一个隐患是业务事务和发消息的非原子性。你在MySQL事务里提交了订单,事务提交后发MQ。如果发MQ失败了怎么办?回滚事务?事务已经提交了。重试?要写重试代码,增加复杂度。如果把发MQ放在事务里,那MQ的IO会成为事务的一部分,极端情况下事务会因为MQ抖动而回滚,业务直接失败,这显然不可接受。
第二个隐患是Binlog方案能捕获所有数据变更,而双写只能覆盖当前代码。说具体点:如果你有一个管理后台可以直连数据库改数据;如果你有一个数据订正脚本批量修数据;如果DBA手动改了一条数据……这些变更,业务双写方案全都感知不到,但Binlog方案能感知到,因为这本质上是MySQL层面的变更流,不是业务层面的调用。这也是我坚定选择Binlog + MQ方案的原因:它把“同步”这件事的起点,从“业务代码”提前到了“数据库日志”。
2.3 MQ选型对比:Kafka和RocketMQ的取舍
在选择了Binlog + MQ之后,另一个决定是选哪个MQ。如果你的公司已经有成熟的中间件团队,通常不用纠结,用公司主推的即可。如果没有历史包袱,我给你一个基于我个人经验的参考。
RocketMQ在顺序消息、事务消息上支持更完善,RocketMQ的“顺序消息”按消息队列粒度保证,同一订单ID的Binlog事件会落同一个队列,天然解决乱序问题;它对消息过滤、Tag的支持也更细,可以在Broker端做轻量过滤,减少无效消费。Kafka的优势主要在大吞吐、生态成熟、流处理一体(Kafka Streams、Flink Connector),但如果要在Kafka里保证顺序,只能依赖单分区,而单分区的吞吐上限会低不少;而且要处理Kafka的Rebalance坑,消费者组里成员变化会触发分区重分配,可能引起重复消费或短时消费停顿。
我的个人建议是:如果数据量巨大(每秒几万以上)且已经有Kafka基础设施,选Kafka但接受局部乱序或用单分区换顺序;如果数据量中等(每秒几千到一两万)且对顺序要求高,RocketMQ会更舒服。我项目里用的是RocketMQ,顺序消息帮我们省掉了至少30%的乱序处理代码。
重要提醒:选MAQ不是终点,选完消息模型后必须立刻处理两个核心问题——消息顺序性和消息幂等性。顺序性问题用消息Key(比如主键ID或业务ID)路由到同一个队列解决;幂等性要靠ES文档里的版本号或更新时间的比对来保证。这两点是后面实操章节的重要铺垫。
3. 异步写ES的工程细节:从理论到可落地的方案
3.1 核心链路设计:Binlog -> Canal -> MQ -> 消费者 -> ES
我直接给出生产环境验证过的链路设计,并标注每一环的作用。
整体链路的完整流程是:
- 应用正常写MySQL,操作产生Binlog。
- Canal伪装成MySQL从库,拉取Binlog并解析成增量事件。
- Canal将增量事件发布到MQ,消息体里包含操作类型、表名、主键ID、变更前后的列数据。
- 消费者从MQ拉取消息,组装成ES文档写入请求。
- ES更新文档到索引,搜索引擎对外提供检索服务。
在这条链路里,每个环节都有可优化、可加固的细节。下面拆开说。
3.2 Canal端的关键配置与使用经验
Canal的部署模式我没必要展开讲,但有几个核心配置值得拿出来说。第一个是canal.instance.master.address,指向MySQL主库地址;第二个是canal.instance.dbUsername和canal.instance.dbPassword,Canal需要专门的账号,这个账号只需要SELECT、REPLICATION SLAVE、REPLICATION CLIENT权限,千万别用生产库的root账号,安全风险太大;第三个是canal.instance.filter,可以用正则过滤表,只订阅需要的表,减少无用的Binlog流量;第四个是canal.mq.topic,决定消息发到哪个Topic。
还有两个实操建议。一是Canal客户端或MQ Producer这边尽量开启批量模式,Canal默认支持按批拉取Binlog日志,本身就有批量能力,一拉就是几百条,下游发布MQ时按批发布,吞吐会好看很多。二是Canal消费BitBinlog的startPosition默认从最新开始,也可以手动指定位点做历史数据回放。
一个我踩过的坑:Canal所在机器和MySQL主库的时钟偏差问题。Canal判断位点时基于Binlog的文件名和偏移量,不依赖时间,所以如果直接部署,影响不大。但如果你在数据订正或者重放逻辑里用到了CanalEvent里的执行时间字段(executeTime),那两台机器时钟不同步时会误导你判断消息新旧。我建议机器统一配NTP时钟同步,或者在消息体里带上MySQL侧的事务提交序号作为逻辑时钟。
3.3 消息体设计与顺序保证
MQ消息不能只丢一坨JSON,必须精心设计字段。我常用的消息体模板如下:
{ "id": 123456, "table": "product", "type": "UPDATE", "data": { "id": 123456, "title": "新版轻量双肩包", "price": 299.00, "stock": 500, "update_time": "2025-04-06 10:30:00" }, "old": { "price": 259.00 }, "ts": 1743903000000 }字段含义如下:id是主键,用于ES文档ID;table是表名,消费者据此区分索引;type是操作类型,INSERT/UPDATE/DELETE;data是变更后的完整行数据;old是变更前的部分字段,主要用于精确判断哪些字段变了,减少无用的ES更新;ts是事件时间戳。
有了完整消息体,消费者就能基于完整行数据重建ES文档。为什么强调传全字段而不是只传增量列?因为如果只传一行更新过的列,ES侧的文档其他字段就丢了或需要你手动做合并,合并逻辑在并发更新场景很难做对。传全字段是成本最低、逻辑最清晰的方案——反正ES文档本来主要就是在索引时保存一份完整数据。
顺序保证这块,做法很直接:用业务主键ID做MQ的消息Key,在RocketMQ里Message的setKeys传入ID。RocketMQ的顺序消息按MessageQueue维度保证,Producer通过MessageQueueSelector根据Key选择队列,同一个ID的消息永远进同一个队列,消费者单线程或者串行消费该队列,这样同一条数据的更新就严格按顺序执行了。在Kafka里等价的方案是选择分区时用Key做Hash,同一个Key落同一个分区,注意分区数变化会导致相同Key落入不同分区,所以顺序性在分区数增减后会被打破,分区数提前规划好,别频繁扩缩容,扩了就无法保证顺序。
3.4 消费者幂等与版本号的硬核处理
顺序问题解决了,还有幂等。消息队列最常见的语义是至少一次(At Least Once),这意味着消费者收到重复消息是常态,特别是在消费超时、Rebalance、手动ACK失败等场景下。如果消费者不处理重复消息,ES文档就会被旧数据覆盖——后果是标题改了又被打回原样。
我推荐的做法是版本号优先。在MySQL表里加一列version(整型),每次更新都SET version = version + 1。Binlog事件里带上这个version。消费者在组装ES文档时,把version存进ES文档的一个字段。写入前先根据文档ID查一下ES里旧文档的version,如果新消息的version <= ES文档的version,直接丢弃(或打日志)。如果ES里文档不存在,直接写入。
这个策略有一个细节要注意:查ES再写ES有一个时间窗口,两个并发消息同时进来可能都查出旧版本,然后都写入,导致后写的覆盖了先写的。为了解决这个,可以在ES里用version字段配合乐观锁条件更新,比如在Update脚本里加条件:if (ctx._source.version < params.newVersion) { ctx._source.title = params.title; ctx._source.version = params.newVersion; }。这样把“版本校验”下沉到ES侧原子操作,就不怕并发了。
// ES Update脚本示例 { "script": { "source": "if (ctx._source.version == null || ctx._source.version < params.newVersion) { ctx._source.title = params.title; ctx._source.price = params.price; ctx._source.version = params.newVersion; }", "lang": "painless", "params": { "newVersion": 3, "title": "新版轻量双肩包", "price": 299.00 } } }这条脚本在生产上实测很稳,是保证最终一致性的核心防线。
3.5 路由策略与ES分片设计
ES索引底层是分片(Shard),文档写入ES时会先计算路由,落到具体分片上。默认路由是_id的哈希。如果你在构建业务查询时用了自定义路由(比如按用户ID查询其下的所有文档),建议写入和查询使用同样的routing策略,否则会“写入在分片A,查询却去分片B找”。
我见过一个坑:某团队在索引创建时指定了routing为userId,查询时确实传了routing=userId,返回正常。但是后台有批量导入脚本直接按_id写入,没有传routing,导致一部分文档被路由到别的分片,查询时永远命不中。排查很久才发现是routing不一致。
所以我的建议是:除非有清晰的按维度查询需求(比如“查某用户的所有订单”),否则默认就用_id路由;如果确实要按业务维度路由,全链路写入和查询都强制传同一个routing字段。另外注意routing字段的选择会影响分片分布是否均匀,如果选userid这种分布比较离散的字段,问题不大;如果选status这样的枚举字段,会造成分片热点,写入全打在某个分片上。选routing字段时先想清楚它的基数够不够高。
4. 一致性兜底机制:光靠同步链路不够,必须加保险
4.1 定时对账任务:最朴素也最有效的兜底
即便Binlog + MQ链路再完善,线上也会出意外:Canal挂了几小时没人发现,MQ积压了几百万条消息消费不过来,或者消费者代码有Bug导致某张表的数据一直没更新。这种场景下,同步链路天然是哑火的——它不会告诉你“我有问题”,它只会安静地让数据不一致。
所以一定要加一层对账任务。原理很质朴:定时任务把MySQL里的数据按主键分批捞出来,和ES里对应文档比对核心字段,不一致就报警或者自动修复。
比如对订单表,每小时跑一次对账任务,按订单ID范围分页拉取MySQL记录的id、status、update_time,再根据id批量查ES文档,比对status是否一致。Mysql查询用WHERE update_time >= ?只捞最近一段时间变更过的数据,减少扫描量;ES用mget或者terms查询一次性拉1000条,减少请求次数。如果发现不一致,把条目记录下来推给告警系统,或者干脆触发修正逻辑重建文档。
这套方案看起来土,但确实救了无数的线上故障。我前东家的搜索团队就是这么干的,每小时对账一次,多次在业务低峰期修复了因为线上变更导致的产线数据不一致。
4.2 增量重放与全量重建机制
对账发现不一致后,需要一个快速修复的手段。增量重放指的是重新把MySQL里的相关行数据读出来,重新组装文档写入ES。实现方式有两种:一种是查MySQL的单行或分页数据,组装好之后直接写ES(update接口);另一种是把需要重放的主键列表推给MQ,让消费者走正常的同步链路去消费修正消息。
另一种更彻底的手段是全量重建。当ES索引的映射需要升级(比如给字段加分词器、改类型)、或者数据错乱太严重、或者索引要重建分片时,就需要全量重建。全量重建的基本思路是重建一个新索引,把MySQL数据全量灌入,构建好之后做别名切换(alias),然后把旧索引删掉。全程要注意“写新索引的同时,增量消息还要继续进来,别丢”,所以全量重建通常配合“先灌存量,再追增量”两阶段方案。增量追平后,把新索引的alias指向业务使用的别名,实现无缝切换。
全量重建我建议用并发分片的方式:把主键范围分成多个区间,每个区间一个线程并发扫描MySQL、组装文档、批量写入ES。我之前用16个并发线程重建过千万级文档的索引,大概十几分钟跑完,比单线程快了一个数量级。注意分批提交大小别太大,ES批量写入建议1000~5000条一提交,太大会导致ES内存压力飙升甚至OOM。
4.3 延迟兜底与最终一致性窗口
即便所有机制都到位,架构上也要接受“最终一致”的现实,而不是“强一致”。ES默认refresh间隔1秒,也就是说文档写入后,最长1秒内才能被搜索到。这个窗口是我们改变不了的物理限制。但我们可以用业务手段规避用户感知:比如商品名称或价格更新后,用户如果搜到旧值,但跳转到详情页看到的是新值,大多数场景是可接受的——关键在于“细节页强一致,搜索页最终一致”的语义划分要清晰地传达给产品和老板。
另外,如果你实在不能接受ES的1秒延迟,可以把ES的index.refresh_interval调成-1(不自动refresh),改为手动refresh。但这相当于自废武功,吞吐会大打折扣,只适合极少量低延迟写入的场景,不推荐常规业务使用。
5. 方案选型与生产实践的决策建议
5.1 按数据规模和团队情况选型的决策树
聊完方案和细节,回到开头的核心问题:到底怎么选?我整理一个决策树,帮你快速定位。
- 数据量小(单表少于百万级)、对实时性要求不高、团队没有专职中间件运维:先上定时任务扫表 + 对账兜底。这个方案成本最低,也能解决大部分问题。不要一上来就上Canal+Kafka全家桶,运维负担和排查复杂度会压垮小团队。
- 数据量中等(百万到千万级)、有基本MQ基础设施、需要秒级实时性:上Canal + MQ + 消费者的标准链路,重点做好消息顺序和幂等。
- 数据量很大(亿级)、多个业务方都需要消费MySQL变更、对搜索实时性要求很高:上Canal + Kafka(或RocketMQ)+ 多个独立消费者组,同时配置好分片数、消费者线程数和独立索引。
团队能力也是个变量。我见过团队里没有专职DBA,也没有人熟悉Canal,硬上这套方案后,Canal挂了没人知道怎么恢复,反而比定时任务更糟。方案选型不能只看技术先进性,还要看团队有没有能力运维它。如果团队对中间件不熟悉,建议从定时任务+对账起步,等人力和能力齐了再演进到Binlog链路。
5.2 几个我踩过的高频“坑位”提醒
长期维护这类系统,有几个高频问题值得重点规避:
第一,Binlog格式必须是ROW。如果MySQL的Binlog格式是STATEMENT或者MIXED,Canal没法准确定位行级变更。检查方式:SHOW VARIABLES LIKE 'binlog_format'。重点是ROW格式才能从Binlog里读出完整的前后镜像数据。
第二,Canal账号权限别给大。只给SELECT、REPLICATION SLAVE、REPLICATION CLIENT三权足矣。如果给了DDL权限,一旦配置泄露,攻击者能直接改表结构,这是重大安全隐患。
第三,消费者逻辑必须做异常隔离。写ES失败了不要一直重试同一个消息,要把失败消息丢进“死信队列”或者本地记录失败明细,走对账重放流程。如果写ES持续失败还傻傻重试,会拖慢整个消费链路,垃圾消息会把队列堵死。
第四,新增大字段或者大表全量同步之前,先评估内存。ES写文档时如果字段过大(比如几MB的文本),会占用大量堆内存和磁盘,建议对大字段做裁剪或者只同步需要检索的子集。
第五,注意Java和时间类型的处理。LocalDateTime序列化到ES时要统一格式,否则ES侧排序会错乱;日期类型在建索引时定义成date,别省这步。
5.3 后续扩展:从MySQL同步ES到更广阔的数据生态
MySQL同步ES的架构一旦跑顺,会沉淀出通用的数据同步能力,它可以复用的场景很多:MySQL同步到ClickHouse做实时分析,MySQL同步到Redis做缓存预热,MySQL同步到HBase做冷热数据分离,甚至MySQL同步到另一个MySQL做读写分离或灾备。链路本质都是一样的:捕获Binlog -> 解析事件 -> 投递消息 -> 对端写入。
这个能力做厚之后,可以抽象成公司内部的“数据管道平台”,业务方在上面自助订阅自己关心的表结构变更,不再需要每个项目重复造一套Canal + MQ的轮子。这也是我这几年的核心职业体验:一个好的数据同步设施,是用一次成本撬动了无数后续业务的复用红利,极其划算。
写在最后:一些实在的建议
回到最核心的那句话:MySQL和ES之间,不存在“强一致”,只有“最终一致”,我们可以用同步链路把不一致的窗口压到秒级,再用兜底机制把它压榨到可接受的极限。明白这一点,你设计系统时就不会追求不切实际的“绝对一致”,而是把功夫花在“如何发现不一致、如何快速修复不一致”上。
如果让我给一个复习清单,重点记住这几条:优先考虑Binlog订阅方案,别用业务双写硬扛;消息里带上版本号和全字段数据,消费者侧做版本校验和幂等;建索引时就把routing策略想清楚,避免后期再改;定时对账任务必须有,它是发现静默故障后的最后一道防线。
如果这篇文章能帮你在设计同步架构时少踩几个坑,我就很满足了。如果你正在纠结选型,或者已经在生产环境遇到了一致性问题,欢迎按上面的思路逐步排查调整。做数据同步没有银弹,只有把每个环节做扎实,系统才能稳定。