MySQL数据同步到Elasticsearch的四种方案选型与实战指南

📅 发布时间:2026/10/6 3:05:33
MySQL数据同步到Elasticsearch的四种方案选型与实战指南
1. 为什么大家都在做 MySQL 到 ES 的同步先说个常见的场景业务库在 MySQL 里表结构设计得很规范事务、关联查询都靠它支撑。但随着数据量上来特别是搜索结果、模糊查询、聚合统计这些需求一多MySQL 就有点力不从心了。这时候 Elasticsearch后面统一叫 ES就派上了用场——倒排索引天生就是干全文搜索的单字段过滤、范围聚合、高亮命中这些活它干得又快又舒服。但问题来了ES 不是个数据库它的定位是搜索和分析引擎不适合当业务的唯一数据源。数据还是得存在 MySQL 里ES 里保存的算一份副本专门给搜索和查询用。既然有了副本就有了同步的需求MySQL 里插入、更新、删除了数据ES 这边得有对应变化。听起来简单真正动手做的时候你会发现同步两个字的水很深。我见过不少团队前期没选好同步方案上线之后出了各种问题数据延迟太大、丢数据、ES 和 MySQL 对不上账、重建索引费了半天劲。说实话这活儿虽然不难但方案选错了后面全是坑。作为在不同项目里做过多次 MySQL 到 ES 同步的工程师我把主流的几种方案认真梳理了一遍把各自的原理、代码路径、成功场景、踩坑点都摊开来讲。文章的定位不只是给个结论用哪个,而是帮你看清楚每种方案背后的同步机制再结合你的数据量、实时性要求、团队技术栈来选择。期望你读完之后可以直接照着方案去落地不用再摸着石头过河。2. 同步方案的底层逻辑全量、增量与实时在我们逐一切入各个方案之前有必要先把同步这件事本身拆开。很多方案之争其实是同步机制之争。2.1 全量同步一次性把数据导入 ES全量同步的意思很简单就是把 MySQL 里的某张表或者多张表从头到尾读一遍写入到 ES。常见做法有两种直接用 SQL 查出所有数据例如SELECT * FROM orders。配合分页或主键游标分批读避免一次性把几百万行数据捞到内存里活活撑爆。全量同步一般用在两个场景第一次对 ES 构建索引或者 ES 数据因为某种原因整体损坏了要做重建。它的优点是逻辑简单、实现容易几分钟就可以写个脚本跑完缺点是如果数据量很大耗时较长且同步过程中如果 MySQL 的数据还在变你可能读到的是不一致的快照。2.2 增量同步只同步发生变化的数据增量同步就高级多了本质上是只处理新增、修改和删除。实现增量同步有两条技术路线基于时间戳/版本号的轮询表里有一个字段记录数据的最后更新时间同步工具定期查询UPDATE_TIME 上次同步点的数据。基于 binlog 的监听解析 MySQL 的 binlog精确感知到每一行数据的变更操作然后按变更操作同步到 ES。增量同步直接决定了系统的实时性——时间戳轮询能做到分钟级延迟binlog 监听可以做到近实时通常毫秒到秒级。2.3 什么是实时真的需要秒级延迟吗我觉得必须先问清楚一个问题你的业务到底需要多快不少准备做同步的朋友张嘴就是我要实时同步。但真实业务里实时是一个很昂贵的词基于 binlog 的方案要引入额外的中间件比如 Canal还要维护消费进程处理各种异常重试逻辑而基于轮询的方案最多就是每 5 秒查一次数据库如果业务上确实容忍几十秒的延迟那它便宜又好维护。我在电商项目里做过一个搜索需求用户搜索商品时要求商品上下架状态 1 分钟内生效。这种业务 1 分钟左右的延迟完全够用完全不需要 binlog。反而是订单对账这种场景要求 ES 里的金额和 MySQL 完全一致延迟还不能大那就要上 binlog 并且把同步做成端到端的可靠机制。所以别上来就追求最先进先想清楚你的搜索场景对时效的容忍底线在哪再选型。否则要么你花了 10 倍的维护成本解决了一个不存在的问题要么上了简单的方案数据老是对不上。接下来我们把方案一个个过一遍你先记一下这几个关键词JDBC 轮询、binlog 监听、离线批量、CDC。后面每套方案本质上都是其中一种或几种思路的组合。3. 方案一基于 JDBC 轮询的定时同步以 Logstash 为例这个方案应该是很多团队第一个接触到的。逻辑非常直白定一个周期用 JDBC 连接 MySQL按某个条件查出数据写入 ES。3.1 Logstash 的 JDBC 插件工作方式如果你不想写代码直接用 Logstash 的 JDBC Input 插件就能搞定。它的大致逻辑是这样的你在配置里告诉它数据库地址、用户名、密码。设置一个schedule比如*/5 * * * * *表示每 5 秒执行一次。你写一条 SQL比如SELECT * FROM products WHERE update_time :last_run。Logstash 会记住每轮执行的查询结果并记录一个最后执行的时间点下一轮执行时把这个时间点作为参数传进去。一个典型的配置长这样input { jdbc { jdbc_driver_class com.mysql.jdbc.Driver jdbc_connection_string jdbc:mysql://localhost:3306/shop?useSSLfalseserverTimezoneAsia/Shanghai jdbc_user root jdbc_password password schedule */5 * * * * * statement SELECT * FROM products WHERE update_time :sql_last_value AND update_time NOW() ORDER BY update_time ASC use_column_value true tracking_column update_time tracking_column_type timestamp } } output { elasticsearch { hosts [http://localhost:9200] index products document_id %{id} action index } }关键点有几个tracking_column指定增量字段Logstash 用它的值作为游标。document_id必须对应 MySQL 的主键 ID这样 ES 里的文档就是幂等更新不会因为同步多次产生重复文档。schedule决定轮询频率一般不要低于 3 秒频率过高会给 MySQL 造成不必要压力。3.2 这个方案最舒服的场景和最容易踩的坑最舒服的场景数据量中等比如几十万到几百万行、表结构相对稳定、搜索实时性要求不严苛且关注点更多是能快速跑起来。再有如果团队成员熟悉 Logstash 但不会写 Java 代码这方案基本就是最优选了。最容易踩的坑第一个坑是update_time的精度问题。MySQL 的 timestamp 默认精度到秒如果你在一条 SQL 更新里让两条数据的update_time相同那 Logstash 用:sql_last_value做增量时下一次查询会漏掉和上次游标相等时间戳的数据因为查询条件通常是不是。我的解决办法是应用层写更新时间时用CURRENT_TIMESTAMP(3)把精度提到毫秒同时尽可能避免批量更新同一张表在完全同一个时间点完成。第二个坑是删除操作。JDBC 轮询本质上只关心行数据的变化它是感知不到DELETE FROM products WHERE id 123的。你查出来的结果里不会有这条数据ES 里残留的就永远删不掉。好一点的补救是业务层的删除改成逻辑删除加一个is_deleted字段同步的 SQL 里带上WHERE is_deleted 0这样删除变成了更新能被同步到 ES。第三个坑是 MySQL 数据库连接长时间空闲导致掉线。Logstash JDBC 插件自带连接池但如果网络环境比较特殊加上 MySQL 的 wait_timeout 设置太短跑一个晚上后插件可能连不上数据库。我在配置里一般会在 JDBC 连接串上额外加几个参数autoReconnecttrueconnectTimeout3000socketTimeout30000。3.3 从工程角度看这套方案的取舍如果你是在小团队、项目初期我其实挺推荐先用 Logstash 跑起来的它零代码、可视化配置、监控也方便。但如果你已经想清楚了未来数据量会快速增长或者要求秒级以内延迟那这套方案可能只是临时的过渡方案——不是它不好而是它应对不了大流量和高实时。这套方案的延展优化空间其实也有你可以自己撸一套 JDBC 轮询脚本不一定非用 Logstash用 Python 的pymysql配合APScheduler也能达到同样效果。但工程上要自己处理游标、异常、重试、监控除非成本真的很高否则没必要重复造轮子。4. 方案二binlog 监听同步以 Canal 客户端消费为例这才是很多大厂在用的实时同步方案。我从原理开始讲不绕弯子。4.1 binlog 是什么为什么它能做到实时binlog 是 MySQL 的二进制日志记录了所有数据变更操作。常见的 binlog 格式有三种STATEMENT、ROW、MIXED。对我们的场景来说必须选择 ROW 格式——它记录的是哪一行变成了什么值而不是执行了什么 SQL 语句这样我们从 binlog 里拿到的是最直接的数据变更前后快照。Canal 是阿里开源的一个中间件它会伪装成 MySQL 的从库向主库发起复制协议请求。主库觉得它是一个正常从机就把 binlog 推给它。Canal 拿到 binlog 之后解析成结构化的数据再推给消费端。链路大概是MySQL binlog - Canal 伪装从库 - 接收并解析 - 推送给 MQ比如 Kafka/RocketMQ- 消费端读取消息 - 写入 ES有人会问Canal 一定需要配合 MQ 吗不一定。Canal 自身支持 TCP 直连、MQ、适配器等模式。但在生产环境我强烈建议中间加一层 MQ。因为一旦直接让 Canal 连消费服务只要消费服务重启或处理不过来Canal 的数据就会积压最终可能反过来拖垮数据库的主从复制压力。加一层 Kafka 这样的消息队列消费端可以自由错峰、回放、重试。4.2 数据模型映射与代码示例我以一个订单表为例假设同步到 ES 的order_index索引里面要包含订单号、用户 ID、金额、状态、创建时间这些字段。Canal 会把 binlog 里的变更转成一个CanalEntry.RowData里面有before和after两个字段分别表示变更前后的数据。一个简化版的消费端伪代码长这样public class OrderSyncConsumer { public void handleMessage(CanalEntry.Entry entry) { if (entry.getEntryType() ! CanalEntry.EntryType.ROWDATA) { return; } CanalEntry.RowChange rowChange CanalEntry.RowChange.parseFrom(entry.getStoreValue()); for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { String tableName entry.getHeader().getTableName(); CanalEntry.EventType eventType rowChange.getEventType(); if (order.equalsIgnoreCase(tableName)) { MapString, Object afterMap toMap(rowData.getAfterColumnsList()); MapString, Object beforeMap toMap(rowData.getBeforeColumnsList()); if (eventType CanalEntry.EventType.INSERT || eventType CanalEntry.EventType.UPDATE) { // 写入或更新 ES esClient.index(order_index, afterMap, afterMap.get(order_id).toString()); } else if (eventType CanalEntry.EventType.DELETE) { // 删除 ES 文档 esClient.delete(order_index, beforeMap.get(order_id).toString()); } } } } }看起来不复杂但有几个细节要特别注意混合变更类型Binlog 里有可能出现插入后立刻更新、更新后又更新这种连续操作消费端如果按顺序处理以最后一笔为准中间过程可以跳过。删除操作的索引位置:delete 事件只能拿到before数据你需要从那里取主键 ID再去 ES 删除对应文档。更新频率过高如果一张表在短时间内被频繁更新每次都同步到 ES 会造成写入压力。建议做一个按主键聚合最新的状态覆盖旧状态的小缓冲操作比如每 100 毫秒 flush 一次。4.3 Canal 方案的高可用与常见故障排查高可用部分Canal 本身支持集群模式多台 Canal Server 可以组成集群同一个 destination 可以由集群里的节点消费。但要注意默认模式是一个分区对应一个 Canal 节点高可用体现在故障转移上而不是负载均衡。再就是消费端必须手动提交 offset。我用 Kafka 的时候是消费完数据并成功写入 ES 之后再提交 offset。否则如果先提交了 offset 再写 ES 失败消息就丢了。这属于至少一次的语义ES 侧用幂等document_id来兜底。常见故障找不到 binlog 文件Canal 记录了自己消费到的 binlog 位置但如果你的 MySQL 执行了FLUSH LOGS或 binlog 过期清理Canal 要读的 binlog 已经被删了。解决办法是设置合理的expire_logs_days推荐 7 天以上并监控 Canal 的延迟。字段类型不匹配MySQL 里的decimal类型在 binlog 里解析出来是字符串比如12.34你直接写进 ES 会导致 ES 把它当成 text 或 keyword。消费端要做类型转换。大事务造成延迟飙升MySQL 里一个大事务比如批量更新 10 万行会产生一条很大的 binlogCanal 解析和下发都会变慢最终导致同步延迟猛增。这不是 bug是物理限制需要从业务侧拆批。4.4 什么时候应该用 binlog 方案数据量大、要求实时性高秒级以内、数据变更频繁。这些场景下轮询方案跟不上只有 binlog 方案能打。它还顺带解决了一个隐藏问题MySQL 的物理删除可以被感知到进而同步删除 ES 里的文档。代价就是引入中间件、增加运维复杂度。团队里至少要有一个人能把 Canal 集群搭起来并且熟悉 Kafka 的消费水位。如果你所在的公司已经有 Canal 和 Kafka 基础设施那这个成本其实没有那么高。5. 方案三离线批量同步以 DataX 为例前面两种方案解决的是持续同步的问题但有一个场景漏掉了ES 里的索引如果因为某种原因要全量重建比如索引 mapping 变了doc 需要换一个方式存储这时候你不能一条条地靠 canal 慢慢补太慢了不如直接全量导一次导完再切到增量。DataX 就是用来干这个的。5.1 DataX 是做什么的怎么把 MySQL 数据写到 ESDataX 是阿里开源的数据同步工具特点是异构数据源之间的离线同步。它的架构里有 reader 和 writer具体到我们的场景MySQL reader 读取表数据Elasticsearch writer 写入 ES 索引。一个 MySQL 到 ES 的 DataX 任务 JSON 大致长这样{ job: { content: [ { reader: { name: mysqlreader, parameter: { username: root, password: password, connection: [ { jdbcUrl: [jdbc:mysql://localhost:3306/shop], querySql: [SELECT id, title, price, create_time FROM products WHERE id ?] } ] } }, writer: { name: elasticsearchwriter, parameter: { endpoint: http://localhost:9200, index: products, type: _doc, batchSize: 500, column: [ {name: id, type: keyword}, {name: title, type: text, analyzer: ik_max_word}, {name: price, type: double}, {name: create_time, type: date} ] } } } ], setting: { speed: { channel: 4 } } } }DataX 的querySql可以写你能想象的任意 SQL这意味着你可以在同步过程中做连表操作——把一张订单表和一张用户表 join 起来写到 ES 的一个索引里。这种能力是 Logstash 很难优雅替代的。5.2 用 DataX 做全量 增量的组合拳在实际项目里我很少把 DataX 单独当作永久同步方案而是把它和方案二binlog组合起来第一次搭建索引用 DataX 全量导入 MySQL 当前全部数据。之后日常同步启动 Canal 增量链路。如果 ES 数据因为误操作损坏了可以把增量链路暂停重新跑一次 DataX 全量再恢复增量链路。这里有一个时序问题全量导入的期间业务往 MySQL 写的新数据怎么办。如果全量跑了一个小时这一小时内 MySQL 的新增数据在全量开始时没有但全量结束后 Canal 增量又只从启动增量链路的时间点开始记录那就产生了空窗期。解决思路是先在 canal 里记录一个起始 binlog 坐标对应时间点然后跑 DataX 全量全量完成后再从那个坐标开始增量消费。这样一个小时内的变更就被补上了。要注意 DataX 全量读的是 MySQL 某一时刻的快照如果你用普通的SELECT而不是SELECT ... FOR UPDATE可能读到中间态数据但配合增量消费的补数最终结果是一致的。5.3 跑大表的性能调优DataX 默认的同步速度比较保守大概在几千条每秒。如果你要一次导入千万级数据这个速度太慢了。我在实际用过之后主要从几个方向调优调大channel数量:channel 是 DataX 的并发度4 个并发和 8 个并发效果差距很大。但不要无脑调大每个 channel 会建立独立的数据库连接太大会把 MySQL 连接数打满。分片键拆分如果你的表主键是自增 ID可以在querySql里按照WHERE id start AND id end拆成多个任务片段或者用 DataX 的splitPk配置让它自动按照主键分片。批量写入 ES 的 batchSizewriter 里batchSize控制每批写入 ES 的文档数量我一般设置在 500 到 1000。太少了网络开销大太多了 ES 的 bulk 队列可能扛不住。跑完大表之后再检查一下 ES 里文档总数和 MySQL 的count(*)是否一致。如果不一致问题通常出现在DataX 报错重试导致的重复写入用document_id幂等解决、或者 MySQL 在读取过程中发生数据变更配合增量链路兜底。6. 方案四基于 CDC 流计算框架以 Flink CDC 为例说到 MySQL 同步 ES还有一个在 Java 技术圈越来越常见的方案Flink CDC。这套方案本质上是把binlog 解析 计算转换 写入 ES合并到一个流计算框架里。6.1 Flink CDC 的同步链路工作方式Flink CDC 是 Flink 社区提供的一组连接器可以像读数据流一样读取数据库的增量日志。你写一段 Flink SQL 或 DataStream API它就会持续从 MySQL 的 binlog 里拿到数据变更事件经过你的处理逻辑后sink 到 ES。一个用 Flink SQL 实现同步的例子CREATE TABLE mysql_orders ( order_id BIGINT PRIMARY KEY, user_id BIGINT, amount DECIMAL(10, 2), status VARCHAR(10), create_time TIMESTAMP(3) ) WITH ( connector mysql-cdc, hostname localhost, port 3306, username root, password password, database-name shop, table-name orders ); CREATE TABLE es_orders ( order_id BIGINT PRIMARY KEY, user_id BIGINT, amount DOUBLE, status VARCHAR(10), create_time TIMESTAMP(3) ) WITH ( connector elasticsearch-7, hosts http://localhost:9200, index orders ); INSERT INTO es_orders SELECT order_id, user_id, CAST(amount AS DOUBLE), status, create_time FROM mysql_orders;这个方案的杀伤力在于你在 SQL 里可以随意做类型转换、过滤、甚至 JOIN 多个表的数据后再写入 ES整个同步逻辑变得和写普通 SQL 一样直观。Flink 本身就提供了恰好一次exactly-once的语义保证配合 ES 的 bulk 写入接口能做到很稳的端到端一致性。6.2 和 Canal Kafka 相比Flink CDC 的差异点没有最好的方案只有合不合适的方案。我把 Flink CDC 和 Canal Kafka 放在一起对比一下维度Canal KafkaFlink CDC学习成本需要理解 Canal 的消息格式、Kafka 的消费模式需要了解 Flink 的编程模型和运行环境数据处理能力数据是原始事件需要自己写消费逻辑做转换在 SQL 层就能完成过滤、JOIN、聚合状态管理需要自己管理 offset 和幂等Flink 检查点机制自动保证状态一致性运维复杂度依赖 Canal 集群和 Kafka运维组件多需要部署一个 Flink 集群或 session 模式适合场景有现成 MQ 基础设施消费侧要接多个系统同步链路长、需要做复杂数据转换如果你的下游不仅是 ES还有 Redis 缓存、数据仓库、其他业务库那 Canal Kafka 这种一条 binlog 广播出去各取所需的架构更灵活。Flink CDC 更适合数据从 MySQL 到 ES 一条线走到头的独立任务。6.3 跑 Flink CDC 经常遇到的三个问题第一个是历史数据初始化。Flink CDC 默认不仅能读增量它会在启动时先做一次全量快照然后再切到增量读取。这个特性称为增量快照。但如果数据量大比如单表几百万行默认的并行快照机制可能会短暂占用较多数据库连接。建议设置一个合适的并行度或者先跑批量的离线同步把历史数据导进去再启动 Flink CDC 只消费增量。这样既能避免全量慢也能减少数据库压力。第二个是 schema 变更。MySQL 的表结构如果发生变化比如新增了一个字段Flink CDC 默认会把变更事件捕获到但你的 SQL 如果还停留在旧字段上就会报错。生产环境里表结构变更最好走一个明确的申请流程同步任务要跟着上线变更。第三个是 ES 索引的 mapping 变更。如果你在 Flink 的 SQL 里改变了字段类型比如把amount从 DECIMAL 变成 DOUBLE那 ES 这边的 mapping 也要同步调整不然写入会失败。ES 里的字段类型一旦建立就不能直接改只能重建索引或者新增一个字段用别的方式处理这个在设计阶段就要想清楚。7. 方案对比与最终选型建议说了这么多是时候给一张收敛的表了。我把主流的四个方案放在一起从实时性、实现成本、数据一致性、维护难度等维度做一个总结。方案实时性实现成本能否同步删除适合数据量核心依赖Logstash JDBC 轮询秒/分级极低需逻辑删除配合中小规模Logstash MySQL 定时任务Canal Kafka 消费毫秒/秒级中高能binlog 识别 DELETE大中规模Canal 集群 Kafka 消费服务DataX 离线同步离线低全量覆盖任意规模DataX 任务调度Flink CDC毫秒/秒级中高能默认支持所有事件大中规模Flink 集群 MySQL binlog有一些维度没法用表格量化我单独说一致性方面JDBC 轮询默认是最终一致但因为拉取间隔和大量并发更新可能造成短暂的不一致状态Canal 和 Flink CDC 如果在消费端做好幂等基本可以做到端到端最终一致但要注意删除事件的处理顺序。DataX 是快照级别的中途不停写数据的话全量出来的结果可能和运行期间发生的新增不一致需要配合增量补数。运维投入方面Logstash 最小一台虚机就够DataX 也不大有调度平台就更好Canal Kafka 最大的成本在 Kafka 集群的维护以及消费端服务的监控Flink CDC 需要维护一个 Flink 运行环境如果你公司本来就有 Flink 集群那另当别论。选型还是回到那三个问题数据量多大实时性要求多高团队能驾驭什么技术栈我自己的经验是这样的小项目、表不多、延迟容忍度高直接上 Logstash别折腾。上了一点量级延迟要求在秒级而团队对 Java 熟悉优先考虑 Canal Kafka。需要频繁重建索引或者要做大表初始化DataX 是 MUST HAVE无论主力方案是什么。团队已经会用 Flink并且下游只有 ES 一个目标那 Flink CDC 可以带来很爽的编码体验和一致性强保证。8. 同步链路中的通用设计要点最后这部分是我做了很多次同步后总结的通用教训。不管选择哪种方案下面这些问题总是要面对。8.1 ES 的 document id 必须等于 MySQL 的主键这事我再三强调。ES 里写入文档时用id来保证幂等同一 id 执行两次 index 操作最后一次覆盖前一次。如果同步工具没有显式指定文档 idES 会自动生成一个随机字符串那你下次同步同样一条 MySQL 数据时ES 里会出现两个内容差不多的文档删除更新都会变得一团糟。无论是 Logstash 的document_id、DataX 的 column 定义、还是手写消费代码里的IndexRequest.id都要显式设置。8.2 同步任务的监控与报警同步链路是基础设施一旦延迟或者失败业务方不会马上感知但影响是实际的搜索不到新数据、统计数字不对。所以至少要有这些监控同步延迟Canal 每次消费到的事件的时间戳和当前时间的差值。消费积压数Kafka 里的 lag 指标或者 Logstash 每次轮询的耗时。失败重试数写入 ES 失败的次数和原因。数据对账可以每天用 count 或 max(id) 对比 MySQL 和 ES 的差异虽然对账一次的成本不小但对核心索引非常值得做。我在之前的项目里就是靠着 daily count 对账发现了一个隐藏的 bug——业务侧有极少量的记录通过一个不走 binlog 的通道做更新导致 ES 数据偏少。没有对账这个问题可能几十天后才会被搜索用户发现。8.3 从 MySQL 到 ES 的字段类型映射是同步质量的核心我见过很多人同步完数据之后发现ES 里 sort 价格字段排序不对才发现 MySQL 的 DECIMAL 类型被当成了 keyword 存入 ES。这属于典型字段类型映射问题。这里给一个常用的映射参考表MySQL 类型ES 字段类型说明BIGINT/INTlong/integer数值范围要留意DECIMAL/FLOAT/DOUBLEdouble尽量不要用 text 存数值避免排序和范围查询失效VARCHAR/TEXTtext keyword 子字段全文搜索用 text精确匹配和聚合用 keywordDATETIME/TIMESTAMPdate注意 ES 默认格式是 UTC你可能要做时区转换JSONnested / objectMySQL 的 JSON 要展开成 ES 的嵌套结构否则只保留原串如果你想保留 MySQL 里的字符串作为 ES 的精确值字段就给它加一个keyword子字段这样同一个字段既能做全文搜索又能做精确过滤成本就是存储空间变大一些。8.4 同步失败后的修复策略同步不可能永远不出错出错了怎么快速复位才是关键。我的常规三件事做法第一定位故障点延迟积压在 Canal还是 Kafka还是消费端存 ES。这个必须靠日志和监控。第二小修小补如果只是部分文档对不上用一条精确 SQL 查出受影响的数据单独补推一次。第三彻底重建如果数据乱得无法救援暂停增量链路跑 DataX 全量重建索引再恢复增量。切记在至少一次的同步模式下重复投递是常态丢失消息才是异常。所以消费端必须幂等ES 文档 id 设置为业务主键这是同步系统的安全带。9. 写在最后同步方案没有银弹我在不同项目里把这几种方案轮番用了个遍最大的体会是方案本身没有绝对的好坏选择时必须回到自家业务的增量数据规模、实时性诉求、可投入的运维成本这三点上综合做决策。Logstash 上手轻松但天花板低适合初期快跑Canal Kafka 是大型系统的常见标配但是运维成本实打实Flink CDC 让人在编码上舒心但该有的集群支撑一点也不能少DataX 更像是工具箱里的重锤平时不一定用但重建全量数据的时候真能救命。我个人在实际操作中的经验是永远不要认为某个方案是一劳永逸的。业务的表结构在变数据量在涨ES 的索引设计也在迭代同步链路每过半年就应该回头重新评估一次。另外推荐一条实用主义路线——先评估量级。如果日增只有几百行你完全不需要 binlog 那套重型武器如果日增几千万行也不要再用定时任务去硬扛。量级定了方案自然就浮出水面了。做同步系统本质是在实时性和可控性之间做权衡。一旦想清楚这一点工具的选择只是顺水推舟罢了。