ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

MySQL与ES数据一致性:从双写到Binlog订阅实践

MySQL与ES数据一致性:从双写到Binlog订阅实践 先交代背景。我所在的小组从2020年开始引入ES做商品搜索早期业务量小MySQL和ES的数据就靠定时任务全量刷每天凌晨跑一次同步脚本出问题第二天早上人工补数。到了2022年单量涨上去全量刷一次要二十多分钟线上隔三差五出现“搜得到但点不开详情”“刚上架的商品搜不到”这种问题。老板一句“怎么保证MySQL和ES的一致性”把后端组逼到了墙角我开始认认真真把整套一致性方案梳理了一遍。这篇文章不是理论分析而是把我踩过的坑、试过的方案、最终沉淀下来的一套打法完整讲清楚。不管你是准备从零搭建MySQLES双写架构还是正在为数据对不上头疼都可以直接参考这里面的设计思路和排坑经验。1. 一致性问题的根源两条数据管道的宿命1.1 为什么一定会产生不一致先说个最直白的事实MySQL和ES是两个完全独立的存储系统一个面向OLTP事务一个面向全文检索。MySQL靠的是行锁、MVCC、B树保证事务与查询ES靠的是倒排索引、分片副本、近实时刷新保证搜索性能。它们之间没有任何原生的同步协议一切数据流转都必须由你自己搭建。只要存在两条独立的写路径就一定存在原子性问题。一个业务操作比如“上架商品”在主库写入一条记录紧接着要把这条记录以文档形态写入ES。这两个动作无法放进同一个本地事务里也无法靠数据库本身的分布式事务协议如XA去覆盖因为ES根本不参与MySQL的事务协议。一旦第二步失败、进程崩溃、网络超时MySQL有数据而ES没有一致性问题就发生了。更麻烦的是MySQL和ES的数据形态天然不同。MySQL里可能是一张商品表加一张SKU表需要join之后才能给ES一个完整的搜索文档ES端可能还要做嵌套类型、父子关系、多字段分词。这些转换逻辑在同步链路里一旦出错就会造出“数据存在但字段丢失”“JSON结构不对导致写入失败”这类隐蔽问题。这引出第一句要记住的话一致性不是某个工具能一键解决的而是从写入路径、消息传递、失败补偿到对账兜底一整套机制组合出来的。1.2 先分清“强一致”和“最终一致”的适用边界很多团队一上来就说“我要保证强一致”这个诉求本身值得商榷。如果你真的要求MySQL和ES在任何时刻都一模一样那么最稳妥的方案其实是不用ES或者把所有查询都打回MySQL用MySQL的全文索引。只要引进了ES这个异步检索链路物理上就无法做到真正的强一致因为ES的refresh默认是1秒一次写入后立刻搜索是搜不到的。所以实际可落地的目标通常是最终一致性核心指标有三个同步延迟从MySQL变更发生到ES可搜索中间隔了多久。正常情况应该是秒级以内。数据丢失率极端情况下有没有数据彻底丢失。这里靠消息持久化和对账兜底。冲突处理正确性多条并发更新到达ES时最终留下的版本是否符合预期。理解了这个边界后面所有方案设计都有了明确目标不是消灭不一致而是把不一致的窗口缩到足够小并且在窗口结束后能自动收敛。2. 主流同步方案排坑对比2.1 同步双写简单粗暴但坑最多同步双写是很多小团队的第一反应在业务代码里先写MySQL成功后立刻调用ES API写入失败就抛异常回滚。听起来简单但坑非常密。第一个坑是性能放大。一次商品上架MySQL写一遍ES写一遍网络RT直接翻倍高峰期ES抖动会反过来拖垮主流程。第二个坑是可用性耦合ES集群哪怕只是短暂熔断商品上架主流程就得跟着报错这在架构上等于把非核心依赖升级成了核心故障源。第三个坑是事务边界被打破MySQL事务已经提交ES写入失败后回滚MySQL是做不到的除非做补偿事务这时候只能把失败数据丢进重试表最终还是要回到异步补偿。我不建议在正式环境直接用同步双写扛写入流量除非你的ES写入极其简单、业务量小到可以忽略ES故障影响。但同步双写的思路有一个可取之处它把“MySQL成功”和“ES成功”的边界暴露得很清楚倒逼你想清楚失败后怎么办这个价值不能忽视。2.2 异步双写 MQ解耦与削峰的标准做法演进方向很自然MySQL写入成功后不直接调ES而是发一条消息到MQKafka/RocketMQ由独立的消费者写入ES。这样写入路径被解耦MySQL主链路只关心发消息ES消费能力不够时可以通过堆积削峰ES短暂不可用也不影响主流程。这个方案的核心难点在消息层面消息丢失直接发MQ如果MQ不可用怎么办本地消息表、事务消息RocketMQ都是常见解法。重复消费消费端必须幂等ES的doc_id用业务主键天然支持幂等覆盖写。顺序问题同一业务主键的多个变更如果被并行消费后到的可能先执行导致ES最终是旧数据。顺序问题是异步双写最容易翻车的点。我在项目里被迫加上了分区有序策略相同业务主键的变更消息用主键hash作为Kafka分区key保证进入同一分区消费者单线程消费该分区。代价是单分区吞吐受限但对非超高并发场景完全够用。还有一点容易忽略事务消息和本地消息表的选型。RocketMQ事务消息实现相对干净但依赖RocketMQ部署本地消息表则是在业务库里建一张outbox表事务提交的同时写一条消息记录由定时任务扫表投递。后者实现笨重但很稳MySQL和MQ之间即使断网消息也不会丢。2.3 Binlog订阅同步业务无侵入的终极解法如果想彻底绕开业务代码直接监听MySQL的binlog把变更事件解析出来同步到ES那就是Binlog订阅方案。主流工具是Canal和Flink CDC。Binlog订阅的优势非常明显业务无侵入不用改一行业务代码不用在Service层塞同步逻辑。事件完整MySQL里所有真实发生的insert/update/delete都会进binlog不存在业务漏发消息的问题——很多异步双写方案的漏发本质就是业务代码里某个分支忘记发消息。顺序天然有序单库单表的binlog按写入顺序排列只要消费者按序处理不会出现后写先到的问题。落地时最大的坑在于binlog格式。必须把binlog_format设成ROW还需要binlog_row_imageFULL否则拿到的是partial imageupdate事件里可能只有主键和被修改字段没有完整行数据ES文档就没法重建。另外用户需要有REPLICATION SLAVE, REPLICATION CLIENT权限这块是DBA配合的事项。Canal的部署架构一般有两种形态Canal Server AdapterServer负责拉取解析binlogAdapter如es6 adapter直接消费并写入ES配置JSON映射规则即可。适合快速搭建。Canal Client模式自己在Java/Scala项目里写一个CanalClient订阅消息后自由处理。适合有复杂转换逻辑、需要多路分发同时写ES和Redis的场景。Flink CDC则更重量级它把binlog解析后的数据流暴露给Flink算子配合checkpoint可以做精确一次语义exactly-once的同步。如果已经上了Flink生态从MySQL同步到ES、ClickHouse、Kafka都可以用一套管线搞定这是它最大的优势。代价是集群运维成本高如果只是为了同步两张表而布一套Flink集群性价比并不高。2.4 定时全量重建永远不要放弃的兜底方案不管同步链路做了多少加固线上总会有一些“说不清”的数据问题某张表DDL变更导致Canal解析失败、某个ES mapping字段类型冲突导致批量写入被拒、某次发布时消费组挂了一下午没人发现……这些情况下增量同步链路没法自愈必须有一个兜底手段。定时全量重建就是干这个的。常见做法是每天凌晨或业务低峰期通过SELECT * FROM table WHERE update_time 上次水位做增量拉取注意依赖update_time字段且要求该字段在更新时可靠刷新或者干脆全表扫一遍按主键批量写入ES利用doc_id覆盖写实现幂等。但是全量重建有个隐蔽的坑全量扫库会影响源库性能。大表没加索引的update_time查询一次全表扫描就能把慢查询打满。我做过的最优解是分页流式查询配合游标每页1000条不做深分页不用offset翻页用主键大于上一页最大ID翻页同时限制读流量避开业务高峰。定时全量重建不能作为主要同步手段但必须作为最后一道安全网。真正的一致性体系增量链路管“快”定时对账管“准”两者缺一不可。3. 最终一致性落地的核心设计3.1 版本号机制让旧数据自动失效异步链路最普遍的问题是乱序或重复投递。比如一个商品先更新了标题又更新了价格两条消息进入消费队列后如果第一条因为网络抖动延迟处理第二条先写入了ES紧接着第一条到达就会把ES里的新价格覆盖成旧价格——典型的数据回退。解决乱序最优雅的方式是版本号机制。在MySQL业务表里维护一个version字段每次更新version version 1同步到ES的文档里存储这个version写入ES时通过_update脚本或version_typeEXTERNAL_GTE的方式只有新版本的version大于等于ES中已有文档的version才执行写入否则丢弃。ES原生的version_typeEXTERNAL_GTE可以很好地支持这个逻辑指定一个外部版本号ES会以它作为写入依据若外部版本号大于等于当前文档版本则覆盖低于则返回409冲突。但需要注意ES的外部版本号是带类型的数值可以直接用业务表的version字段映射过去省去自己写脚本判断的麻烦。没有version字段的老表也可以用update_time来代替。ES文档里存时间戳消费端更新前先比对时间戳旧数据不写入。不过时间戳精度问题要注意如果MySQL的datetime只精确到秒同一秒内的两次更新可能无法区分先后建议至少在ES端存储毫秒级时间戳或者在特殊场景下升级到version字段方案。3.2 幂等写入把重复消费的伤害降到零MQ消费最常见的语义是at-least-once消息可能重复投递所以写入ES的操作必须幂等。好消息是ES根据doc_id覆盖写天然幂等同一个doc_id执行多次index写入结果和最后一次写入一样不会产生重复文档。但引入了_update脚本之后幂等就需要自己保证。比如“给商品库存减一”这种操作如果写成script里的ctx._source.stock - 1重复消费两次就会把库存减两次产生严重事故。正确做法是所有ES侧的更新都必须是覆盖式写入而不是增量修改。增量修改的逻辑应该在MySQL侧完成ES端只是把MySQL算好的最终结果同步过来。我在这里吃过一次亏。当时为了省流量把“新增评论数”设计成es script自增结果MQ重试时评论数翻倍。最终把所有script改造为全量覆盖风险等级立刻降下来了。注意凡是使用ctx._source.xxx 1这类脚本请反复确认消费幂等性。覆盖式写入虽然流量大但安全和可控压倒一切。3.3 重试与退避别让失败像滚雪球一样扩大ES写入失败无法完全避免集群内存熔断、分片迁移、422 mapping冲突、网络超时……失败消息不能直接丢弃否则数据就悄悄丢了。重试机制分两层第一层同步重试。针对网络抖动、429限流这类瞬时错误可以在消费线程内直接重试2~3次间隔用指数退避比如1s、2s、4s。注意ES客户端的retry_on_conflict参数只能处理版本冲突重试处理不了HTTP层面的连接超时网络重试逻辑要自己写。第二层异步补偿。同步重试仍然失败的消息要投递到死信队列或落库如es_sync_retry表由独立定时任务周期性扫描按失败次数递增退避。同时记录失败原因方便人工介入。这里有一个容易踩的坑不要把大量失败消息全部原样重投到同一个MQ主题。当ES集群故障恢复时所有堆积的失败消息一次性涌入直接把ES打挂这就是重试风暴。正确做法是给死信队列做限流每次只放行N条且按业务主键去重保证同一条数据不会同时被多个任务重复补偿。3.4 对账补偿最后的防线增量链路再完美也没有人能保证万无一失。对账是发现无声失败的最后手段具体做法是定期每天或每小时对比MySQL和ES两份数据的不一致情况。对账设计要避免全量对比否则数据量大到根本跑不完。我推荐增量对账 抽样对账结合增量对账以MySQL的update_time为水位线拉取最近N小时内变更过的记录主键然后去ES里按主键批量mget比对关键字段标题、价格、状态、version。不一致的记录进补偿队列。抽样对账对存量数据按主键hash抽样如千分之一全字段比对用于发现历史遗留问题。对账不是同步链路的一部分而是监控链路的一部分。它的产出应该是告警和补偿任务而不是实时修复。很多团队把对账做成实时同步这是本末倒置对账的价值在“发现系统性问题”比如某天binlog被误删、mapping字段未升级都是靠对账发现的而不是靠重试支撑起来的。4. 实践中容易踩的坑与排查实录4.1 物理删除是同步链路最大的天敌Canal和MQ方案都默认订阅的是“变更”事件。MySQL里如果执行了DELETE FROM product WHERE id1binlog里确实有删除事件但Canal解析出来后ES端的处理需要额外注意——如果直接删除ES文档那么所有基于该文档的搜索、聚合、推荐链路都会瞬间失去它。更隐蔽的问题是很多业务不关心删除事件只订阅insert/update或者解析工具配置没写delete规则于是ES里永远留着MySQL已删除的文档用户搜到一个已下架商品点进去404体验极差。我的经验是尽量用逻辑删除替代物理删除is_deleted字段置1同步到ES后通过filter过滤。如果非要物理删除则必须在同步管道里明确声明delete处理逻辑并对delete事件单独告警跟踪防止漏删。同时要注意ES文档删除之后如果又有幂等重放文档会复活且带旧字段这种诡异问题很难排查。4.2 ES routing配置错乱引发的数据“幽灵”ES的分片路由默认靠_id的hash决定文档落在哪个主分片。如果业务设置了自定义routing比如按用户ID路由那么查询、更新、删除都必须在请求里带上相同的routing值否则请求会被路由到错误分片出现“数据明明存在却查不到”的幽灵现象。我在项目里就翻过车同步管道写入数据时指定了routinguser_id但线上查询接口漏传routing结果搜索端用默认路由去查返回的文档集是另一个分片上的数据整个检索结果错位。排查了很久才发现是routing不一致造成的。更麻烦的是routing值变化的场景比如一个订单文档原来按buyer_id路由后来因为业务扩展改成按seller_id路由如果直接更新文档但routing变了ES会先在原分片删掉旧文档再写入新分片——这个先删后写的过程不是原子的中间状态查询会短暂丢数据而且如果新分片写入失败旧文档又已经被删了就会造成永久丢失。所以routing方案确定后不要轻易更改如果必须改请走重建索引流程不要在线更新。4.3 慢查询导致的超时与重试风暴ES集群偶尔变慢是正常现象但如果写入请求每次都超时消费端还在拼命重试就会形成恶性循环ES越慢重试越多重试越多ES越慢。我见过最夸张的一次是同步任务把ES读库CPU打满正常查询RT从10ms飙到2s。排查ES变慢要看三方面分片健康是否有未分配的副本分片大量unassigned会导致集群yellow甚至red。堆内存写入线程池队列是否堆积thread_pool.write.size和queue_size是否被打满。磁盘IO与合并大量写入触发频繁merge会造成磁盘IO毛刺。可以通过调整索引的refresh_interval比如从默认1s改成30s来缓解高频写入压力。核心原则是同步写入必须与查询链路隔离。如果条件允许给同步管道建一个专用的写索引alias或者错峰写入。这也是为什么异步MQ方案比同步双写更适合生产环境——它天然把ES压力隔离到消费端还能通过调整消费并发来适配ES的处理能力。4.4 bulk批量写入大不等于快同步管道为了提升吞吐往往会批量攒一批数据再写入ES。这个思路本身没问题但要注意batch size选择的“两个极端”batch太小比如10条请求数量太多吞吐上不去batch太大比如5000条单次请求体过大ES需要构建大文档数组容易触发内存熔断反而拖垮集群。我实测下来每批500~1000条单条文档控制在10KB以内每秒总吞吐控制在ES集群可承受范围的80%以内是安全和性能之间的甜点。不要盲目加并发要看ES的写入线程池在请求体到达时有没有排队排队长了说明ES处理不过来加并发只会加剧压力。如果用了BulkProcessor还要注意它的flush_interval和concurrent_requests参数前者控制攒批时间后者控制允许同时在途的请求数。设置不当会导致批量请求乱序前一批还没写完后一批已经开始覆盖对有版本号机制的场景影响不大但没有版本号兜底时乱序可能导致旧覆盖新。5. 工具选型参考与经验总结5.1 Canal、Flink CDC、自研同步器怎么选不同规模、不同团队基因选型差别很大。我整理一张对照表方便直接做决策维度CanalFlink CDC自研同步MQ消费业务侵入无binlog订阅无binlog订阅有需发消息实时性秒级秒级可精确一次取决于MQ消费速度复杂度中等轻量高需Flink集群中等需MQ基础设施适合场景Java技术栈为主、表数量可控已上Flink生态、多端同步需要业务内做转换、规避数据权限顺序保证单表有序依靠算子链/分区需按主键分区如果团队已有Flink环境且数据同步链路不止ES一个目标我会推荐Flink CDC因为它能把ES、ClickHouse、Kafka等一次性打通后续扩展成本低。如果没有Flink环境并且同步目标只有ESCanal自研消费端是性价比最高的方案轻量、稳定、可观测。自研同步器适合同时需要结合业务逻辑做转换、且不想依赖外部同步组件的情况但要自己扛住消息重复、乱序、失败补偿这些问题投入不低。5.2 数据一致性问题的排查方法当线上已经出现“ES和MySQL不一致”的投诉不要慌按这个顺序排查先确认方向是MySQL有数据但ES没有还是ES有数据但MySQL没有判断决定排查重点。前者查增量管道后者查删除/过期逻辑。查消费位点MQ消费组的lag是多少lag持续增长说明消费能力不足lag为0但数据还是不对说明是乱序或转换逻辑问题。查最近一次失败错误日志里有没有批量写入失败死信队列里堆了多少消息手动对单条选一条已知不一致的数据跑一遍同步管道观察中间每一步的数据形态看是不是字段映射问题。查全量兜底上一次定时全量重建是什么时候全量重建后有没有立刻出现新的不一致这些排查每一步都能落到日志和指标上前提是你的同步管道从一开始就要埋好日志和监控包括消费lag、写入成功率、ES集群健康度、死信队列深度。没有可观测性的同步链路出了故障基本等于盲人摸象再强的排错经验也派不上用场。5.3 我个人的落地建议做了这么多方案对比和坑点梳理最后给你一套可以直接抄的落地建议如果是新项目直接上Canal订阅binlog消费端用Java写一个轻量服务写入ES时开启version_typeEXTERNAL_GTE做版本保护目标索引单独配置refresh_interval为10s左右辅以每小时的增量对账任务兜底。这套组合支撑日百万级数据变更绰绰有余。如果是已有老项目、代码里已经塞满了同步逻辑别急着全部推倒重来——先保留现有双写逻辑加一份数据库binlog订阅做旁路校验跑两周对账把双写漏掉的数据捞出来补偿等验证稳定之后再说是否切换到纯binlog方案。渐进式迁移比一步到位风险小得多这在一致性治理里也同样适用。最后再分享一个小技巧所有同步管道的消费逻辑都要支持在测试环境用真实binlog replay。我每次改了ES mapping或同步逻辑都会先拉一份生产的binlog片段在测试环境重放观察最终写入ES的文档是否符合预期。这个习惯帮我提前发现过好几次问题比上线后被线上数据打脸舒服太多了。
返回列表