ARTICLE DETAIL

资讯详情

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

基于Canal实现MySQL到Elasticsearch实时增量数据同步实战

基于Canal实现MySQL到Elasticsearch实时增量数据同步实战 做 MySQL 到 ES 的数据同步最烦的就是双写和定时任务两条路都让人头大。双写侵入业务代码分布式事务难搞失败了对账对到哭定时任务全量扫表时效性差不说还会给数据库平添压力。我最早接手这个需求时第一反应也是能不能基于 binlog 做增量订阅后来在 Linux 环境下用Canal把MySQL 数据实时同步到 ES整个链路跑通一次性解决了列表页搜索、详情页缓存、统计报表等多个场景的数据一致性问题。这篇博文就把整个实战过程完整捋一遍覆盖 Canal 的原理、MySQL 侧配置、Canal Server 部署、消费端写入 ES 以及我实测下来的一堆坑。适合已经会用 ES 但没做过数据同步的同学也适合想把 Canal 从零跑通的读者参考。1. 为什么是 Canal增量数据同步的整体逻辑与方案对比1.1 双写和定时任务为什么不好使先聊聊大家最常走的弯路。业务里最直接的方案就是在 Service 层写代码先写 MySQL再调 ES 接口写入。这个方案的问题在于强耦合订单服务里塞了一堆 ES 逻辑一旦 ES 抖动主流程跟着受影响。就算你引入了消息队列削峰填谷消息丢失、重复消费、顺序问题还是会找上门排查一圈下来你会发现大部分时间都耗在了对账而不是写代码上。另一种常见做法是定时任务全量刷新。半小时跑一次每次SELECT * FROM table WHERE update_time ?。这个方案在数据量小的时候能用但一旦单表过千万分页深翻、大字段传输、ES 大批量写入任何一个环节都会拖垮性能。更关键的是它始终解决不了实时性问题用户刚下单列表页却要等半小时才能看到这在电商场景基本不可接受。1.2 Canal 的工作原理伪装成 MySQL 从库Canal 的核心思路很巧妙它不是去查 MySQL而是去订阅 MySQL 的 binlog。它启动后会把自己伪装成一个 MySQL 的 slave 节点向 master 发送 dump 协议请求。MySQL 主库一看来了个从库就正常地把 binlog 推送给它。这个机制的好处在于对业务代码零侵入MySQL 把 Canal 当从库业务方根本感知不到数据是推送出来的实时性比定时任务高好几个数量级基于 binlog 的 Row 模式可以精确拿到每一行的变更前和变更后数据整个链路的角色划分是这样的MySQL master (开启binlog) ↓ dump 协议 Canal Server (解析binlog转换成事件流) ↓ TCP / MQ Canal Client / Adapter (消费增量事件) ↓ bulk 写入 Elasticsearch我第一次看这张图时最大的感触是Canal 把数据库变更这件事变成了一个标准的事件流下游无论是同步到 ES、Redis、数仓还是触发业务逻辑都变成了订阅消费的问题架构一下子清爽了很多。1.3 完整链路和数据流向实际应用中一个典型的订单同步链路是这样的请求写入 MySQLorders表binlog 记录变化Canal Server 解析出 INSERT 事件消费端拿到行数据后按主键order_id映射到 ES 文档 ID做 upsert 操作。整个过程从提交事务到 ES 可见实测在局域网内大概是几十到几百毫秒的级别。这里要提前建立一个认知Canal 同步的是增量从 Canal 启动那一刻开始产生的事件才能被捕获。历史存量数据需要另做一次全量初始化通常用 Logstash 或 DataX 搞定。两者结合先全量灌入 ES再增量实时追平这个模式在行业内已经很成熟了。2. 环境准备版本匹配与基础软件安装2.1 版本组合建议版本匹配是上手 Canal 最容易踩雷的地方尤其是 MySQL 8.0 和 ES 8 的兼容性问题不少人卡在这里反复折腾。我实测下来比较稳的组合如下组件推荐版本说明LinuxCentOS 7.x / Ubuntu 20.04内核无特殊要求JDK1.8首选/ 11Canal 1.1.5 用 1.8 最稳MySQL5.7.x 或 8.0.x8.0 需要额外处理认证插件Canal1.1.5 / 1.1.7老牌稳定版Elasticsearch7.x推荐 7.138.x 慎用映射配置差异大Canal Adapter与 Canal Server 版本对应如果走 Adapter 路线这个表格里的每个选择都是有原因的。Canal 1.1.5 是我用得最多的版本文档齐全、社区踩坑记录多资料好找到 1.1.7 增加了不少 MQ 相关的支持但核心逻辑变化不大。ES 我强烈建议先用 7.x 跑通因为 8.x 默认开启了安全认证RestClient 的配置方式不一样Adapter 的兼容性也没跟上没必要在第一步就给自己上难度。2.2 JDK 安装要点JDK 的安装本身不复杂但版本别选错。Canal Server 官方要求 JDK 1.8 以上实测在 JDK 11 下也能跑不过有些老版本的 Adapter 在 JDK 11 下会遇到模块化导致的反射报错所以服务器上同时装了多版本 JDK 的话记得把JAVA_HOME指到 1.8。安装完成后三步验证java -version # 期望输出包含 1.8.0_xxx 或者准确的版本号 echo $JAVA_HOME # 期望输出类似 /usr/local/jdk1.8.0_202 which java # 确认 PATH 里用的是你装的那个2.3 MySQL 与 ES 的安装注意事项MySQL 这里我重点提醒两个点。第一binlog 默认是关闭的如果你装的是发行版 MySQL必须手动修改配置文件并重启这一步后面专门讲。第二MySQL 8.0 默认的认证插件是caching_sha2_passwordCanal 连接时可能报认证失败稳妥做法是给 Canal 单独建账号并指定mysql_native_password加密方式。ES 安装有三个容易忽视的细节不能用 root 用户直接启动必须建独立用户vm.max_map_count要调大到 262144 以上JVM 堆内存建议设置为物理内存的一半且不超过 31GB。ES 不像 MySQL 那样需要开启特别的 binlog 之类功能但它对系统参数的要求反而更苛刻启动前把sysctl配置改好能省掉后面一堆莫名其妙的报错。3. MySQL 侧配置开启 binlog 并创建同步账号3.1 binlog 配置参数MySQL 这边是整套链路的源头binlog 没开或者格式不对后面的工作全是白费。编辑 MySQL 配置文件CentOS 通常在/etc/my.cnf[mysqld] server-id1 log-binmysql-bin binlog_formatROW binlog_row_imageFULL expire_logs_days7 max_binlog_size128M逐个解释一下这些参数的作用server-id1主库的唯一标识Canal 作为伪从库的时候会拿它做复制上下文log-binmysql-bin开启 binlog文件名前缀binlog_formatROWRow 模式记录每行数据的实际变更Canal 必须依赖它binlog_row_imageFULL记录变更行的全部字段默认值就是 FULL但显式写出来更保险expire_logs_days7自动清理 7 天前的 binlog防止磁盘被撑爆max_binlog_size128M单个 binlog 文件大小上限到了就滚动生成下一个改完重启 MySQL然后登录验证SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format;两个结果分别应该是ON和ROW。有一点务必注意如果你的 MySQL 本身就属于主从复制架构server-id 千万不能和已有的从库重复。我遇到过一例Canal 起不来日志里反复报 binlog dump 被中断排查到最后发现是 server-id 冲突导致的。每个从库包括 Canal 这个伪从库在复制拓扑里都必须是唯一 IDMySQL 对重复 server-id 的处理非常保守会直接断掉复制连接。3.2 binlog_formatROW 的含义很多人对 Row 模式的理解就是记录变更前后数据但实际操作中它有个重要特性每个变更事件会携带该行的完整字段值。比如一条 UPDATEbinlog 里既有变更前的镜像before image也有变更后的镜像after image。Canal 解析后会把 after 的数据以 JSON 形式暴露给消费者你这个 JSON 几乎可以直接拿来写入 ES 文档。如果用的是 Statement 模式binlog 记录的是 SQL 语句本身Canal 没法还原出具体行数据同步自然无从谈起。这也是为什么 Canal 的官方文档和所有教程都会强调必须用 Row 模式。3.3 创建 Canal 专用账号不建议用 root 账号直接连 Canal职责分离是基本操作。在 MySQL 里执行CREATE USER canal% IDENTIFIED WITH mysql_native_password BY Canal123; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; FLUSH PRIVILEGES;三个权限的含义分别是SELECTCanal 启动时需要读取表结构、元数据等信息REPLICATION SLAVE用于从库复制Canal 正是靠这个权限请求 binlogREPLICATION CLIENT用于获取主库的 binlog 位置信息如果你用的是 MySQL 8.0IDENTIFIED WITH mysql_native_password这段必不可少。如果省略MySQL 会用默认的caching_sha2_passwordCanal 部分版本会报Access denied或者Public Key Retrieval is not allowed之类的错误。5.7 及以下版本不用画蛇添足直接IDENTIFIED BY Canal123就行。4. Canal Server 部署核心配置与启动验证4.1 下载解压Canal 的发行包是直接解压即用的。到 GitHub Releases 页面下载canal.deployer-1.1.5.tar.gz上传到服务器后mkdir -p /opt/canal tar -zxvf canal.deployer-1.1.5.tar.gz -C /opt/canal cd /opt/canal解压后的目录结构里最关键的是两个conf/canal.propertiesCanal Server 的全局配置端口、ZooKeeper、目的地列表都在这里conf/example/instance.properties单个数据同步任务destination的配置管的是连哪个 MySQL、订阅哪些表、消费起点在哪理解这两层配置很重要。一个 Canal Server 可以管理多个 destination每个 destination 对应一个独立的同步任务。默认安装只有一个example任务你可以复制出多个目录来配置不同的 MySQL 实例或者不同的订阅规则。4.2 canal.properties 核心参数打开conf/canal.properties重点关注这几个# 对外提供 TCP 服务的端口客户端连接用 canal.port11111 # 如果用了 ZooKeeper 做集群这里填 ZK 地址 canal.zkServers127.0.0.1:2181 # 当前 Server 管理的所有 destination 列表多个用逗号分隔 canal.destinationsexample # 是否开启 lazy 模式true 表示客户端连接时才加载 destination canal.instance.global.lazytrue单机部署的情况下canal.zkServers可以直接留空但如果有多个 Canal Server 做热备或者对高可用有要求就必须引入 ZooKeeper 集群。这里我不展开集群的搭建细节因为大部分团队的起步阶段都是单机跑通再说。canal.port是消费端连接的端口对应到 Java 客户端里就是CanalConnectors.newSingleConnector(127.0.0.1, 11111, example, , )里的那个 11111。4.3 instance.properties 核心参数conf/example/instance.properties是真正配数据库连接的地方# MySQL 主库地址 canal.instance.master.address127.0.0.1:3306 # 连 MySQL 用的账号密码就是前面建的那个 canal.instance.dbUsernamecanal canal.instance.dbPasswordCanal123 # 从库 ID注意不能和拓扑里其他节点重复 canal.instance.mysql.slaveId1234 # 连接 MySQL 的字符集建议和库表保持一致 canal.instance.connectionCharsetUTF-8 # 订阅规则正则表达式中 \\..* 是转义后的 \..* # 表示默认同步所有库所有表也可以只订阅某张表比如 # canal.instance.filter.regexdbname.tbname canal.instance.filter.regex.*\\..*还有一个经常被忽略的参数在canal.properties里canal.instance.filter.black.regexmysql\\.slave_.*这是默认的黑名单用来过滤掉 MySQL 内部的系统表一般不需要动。如果你发现 Canal 一直在解析mysql库的内部表变更消息很可能就是黑名单配置被清掉了。订阅规则的正则格式是库名.表名注意点和坑也不少匹配所有库所有表.*\\..*匹配指定库的所有表dbname\\..*匹配指定库的指定表dbname\\.tbname多表用逗号隔开dbname\\.tb1,dbname\\.tb2不支持的写法直接写dbname.tbname不带转义会解析异常4.4 启动与日志排查启动命令很简单/bin/startup.sh但启动成功不等于配置成功第一件事是看日志。Canal 的日志分为两块/opt/canal/logs/canal/canal.logServer 主日志/opt/canal/logs/example/example.logdestination 任务日志正常启动时example.log里能看到类似这样的一行INFO canal instance ... start successful同时会打印出它读取到的 binlog 位点信息类似binaryId : 4, position : 123456, scn : 1700000000看到这个就说明 Canal 已经成功连上 MySQL 并开始监听 binlog 了。如果启动报错常见的有三种第一连接不上 MySQL检查master.address的 IP 端口和账号权限第二认证失败多半是 8.0 的密码策略问题回到 3.3 节确认建账号语句第三canal.log里报 ZooKeeper 相关异常如果没配集群把canal.zkServers清空再试。5. 消费端接入Canal Adapter 与自定义客户端的取舍5.1 三种消费方式对比Canal Server 起来只是第一步真正的业务逻辑在消费端。消费方式大体有三条路线方式适用场景优点缺点Canal Adapter只想尽快把 MySQL 数据同步进 ES不想写代码零代码配置即用定制能力弱复杂转换逻辑不好写自定义 Java 客户端需要做字段转换、多表关联、路由规则等灵活可控逻辑全在自己手里要写代码要自己处理 ES 批量写入先推 MQ 再消费多条链路都要消费同一份 binlog 事件解耦多个下游各取所需多一个 MQ 组件运维成本上升实际项目里我见过最多的是第一种和第三种结合先用 Adapter 快速跑通后续发现要加字段映射和清洗逻辑再演进成 MQ 模式。如果你的场景只有把 MySQL 的行数据同步到 ES 的索引Adapter 是性价比最高的选择。5.2 Canal Adapter 接入 ES 的配置Canal Adapter 是官方提供的同步工具它连接 Canal Server 取增量数据然后按配置把数据写入目标端。下载canal.adapter-1.1.5.tar.gz后配置分两层。第一层是conf/application.yml声明 ES 的连接信息和其他数据源信息。ES 相关的一段大概长这样spring: jackson: date-format: yyyy-MM-dd HH:mm:ss time-zone: Asia/Shanghai canal.adapter: conf: canal: server: 127.0.0.1:11111 destination: example src-data-sources: - key: mysql01 type: mysql url: jdbc:mysql://127.0.0.1:3306/dbname?useUnicodetruecharacterEncodingUTF-8 username: canal password: Canal123 es: hosts: 127.0.0.1:9200第二层是conf/es/目录下的映射文件例如orders.ymldataSourceKey: mysql01 destination: example esMapping: _index: orders _id: id sql: SELECT id, order_no, user_id, amount, status, create_time FROM orders upsert: true映射文件里有三个关键点_indexES 索引名需要提前创建好或者允许 Adapter 自动创建_idES 文档 ID 从哪一列取强烈建议用 MySQL 主键保证幂等更新sql读的是 MySQL 的全量查询语句还是增量 JOIN 查询Adapter 依赖这条 SQL 的结果集来构建 ES 文档启动 Adapter/bin/startup.sh观察logs/adapter/adapter.log出现类似sync es 1000 rows或者定时轮询日志说明同步链路已经通了。验证方式很简单改一条 MySQL 数据然后直接curl查 ES 文档看值是否变化。5.3 自定义客户端核心代码逻辑如果 Adapter 满足不了你比如要做多张表 JOIN 之后的结果同步到 ES就需要写自定义客户端。核心代码其实不长关键是理解它的消息确认机制。CanalConnector connector CanalConnectors.newSingleConnector( new InetSocketAddress(127.0.0.1, 11111), example, , ); connector.connect(); connector.subscribe(dbname\\..*); while (running) { Message message connector.getWithoutAck(100); long batchId message.getId(); if (batchId -1 || message.getEntries().isEmpty()) { Thread.sleep(1000); continue; } ListCanalEntry.Entry entries message.getEntries(); for (CanalEntry.Entry entry : entries) { if (entry.getEntryType() ! CanalEntry.EntryType.ROWDATA) { continue; } CanalEntry.RowChange rowChange CanalEntry.RowChange.parseFrom( entry.getStoreValue()); String tableName entry.getHeader().getTableName(); // 根据 rowChange.getEventType() 判断 INSERT / UPDATE / DELETE // 遍历 rowChange.getRowDatasList() 获取每行数据 // 将 afterColumnsList 转成 Map组装成 ES 文档 } // 确认这批次已消费Canal Server 才会删除对应数据 connector.ack(batchId); } connector.disconnect();这段代码有几个地方必须说透。为什么用getWithoutAck而不是get这涉及 Canal 的 ACK 机制。getWithoutAck表示取数据但不确认你拿到数据后可以做任何处理处理成功后再ack(batchId)Canal Server 确认这批次消费完成并删除缓存。如果处理失败或者进程崩溃你可以不 ack甚至主动rollback(batchId)这样 Canal Server 会在下次连接时重新推送这批数据。这个机制和 Kafka 的提交 offset 类似保证的是数据不丢。为什么消费完了才 ack如果你调用get时立刻自动 ack万一写入 ES 失败这条数据就永久丢了。先取后 ack配合失败告警 重试机制才能保证两端最终一致。批大小的选择getWithoutAck(100)的 100 是批次大小。批次太小吞吐上不去批次太大单批处理时间过长遇到断连时回滚的量也大。我实测 500 左右在单表订单场景下比较均衡具体需要压测。拿到RowChange后ES 侧直接用的是 Bulk APIBulkRequest request new BulkRequest(); // 根据 eventType 构造 IndexRequest / UpdateRequest / DeleteRequest // DELETE 事件直接按文档 ID 删除不用带数据 // INSERT / UPDATE 统一处理成 upsert按主键映射 _id client.bulk(request, RequestOptions.DEFAULT);5.4 ES 端索引与 ID 映射不管是 Adapter 还是自定义客户端ES 索引的设计都在最前面。我的建议是文档_id直接用 MySQL 主键这是实现幂等的基础时间字段在 ES 里统一用date类型并指定format避免默认格式解析出错金额字段不要用float或double精度会丢用scaled_float或 long 存储分需要分词搜索的字段用textik分词器仅排序或精确匹配的字段用keyword这里额外提一个典型问题MySQL 的datetime和 ES 的date类型默认格式不一样MySQL 是yyyy-MM-dd HH:mm:ssES 默认是带毫秒的yyyy-MM-ddTHH:mm:ss.SSSZ。如果两边格式不对齐同步进去的值可能解析失败或者不符合预期。字段映射里显式声明format: yyyy-MM-dd HH:mm:ss||yyyy-MM-ddTHH:mm:ss.SSSZ||epoch_millis能兼容多种格式输入是稳妥的做法。6. 实战踩坑数据对不上、写不进去、同步中断三类问题6.1 启停之后消费位点丢失这是我最开始踩过的一个坑Canal Server 跑得好好的重启一次后自定义客户端发现从重启后的新数据开始消费重启期间的增量数据全丢了。原因是 Canal 的位点position记录在 meta 文件里默认存储在conf/example/meta.dat。正常关闭时Canal 会把当前消费到的 binlog 位置写进去但如果是kill -9强杀进程或者 ZooKeeper 模式下位点管理出现了异常重启后 Canal 会重新向主库询问当前最新位点从最新开始监听之前那段时间的数据就漏掉了。解决思路有两个。第一部署时做好优雅停机不要随意杀进程。第二也是最关键的不要指望 Canal 位点永远不会丢一定要在消费端做好数据比对和补偿机制。最简单实用的做法是每隔一段时间跑一次 count 对比发现不一致再触发全量重新同步。Canal 是尽量不丢而不是绝对不丢这个认知得放在心里。6.2 时间字段差了 8 个小时同步到 ES 的时间字段比 MySQL 少了 8 小时这个问题几乎人人都会遇到。根因是时区Canal Server 的 JVM 默认时区是 UTC解析 binlog 里的datetime字符串时按 UTC 转换之后就少了 8 小时。解决办法是在启动脚本里显式指定时区。编辑/opt/canal/bin/startup.sh在 JVM 参数里加一行-Duser.timezoneAsia/Shanghai同理Canal Adapter 的application.yml里也要设置time-zone: Asia/Shanghai。改完重启验证一条数据的时间字段确保和 MySQL 一致再继续。这个坑之所以隐蔽是因为它不一定每次都报错很多时候 ES 里能写入但查询结果就是不对。等到业务方来质问为什么订单时间对不上时排查方向已经绕了一大圈。所以建议配置阶段直接就把时区固定住不要留到后面。6.3 类型转换与写入冲突binlog 解析出来的字段都是字符串直接塞进 ES 时经常会出类型冲突。典型的例子MySQL 的DECIMAL(10,2)解析出来是字符串123.45ES 映射是scaled_float如果 mapping 里没指定合适的格式写入报错MySQL 的INT解析出来是字符串123ES 映射是integer默认能自动转换但某些版本下bigint超出 JS 安全整数范围会被截断NULL字段在 JSON 里可能被序列化成nullES 的字段映射如果禁止了 null写入会失败应对方式就是在消费端做一层显式类型转换。我自己习惯在自定义客户端里写一个convertByEsMapping方法根据目标 ES 字段类型逐个转换把字符串转成对应的 long、double 或日期类型。用 Adapter 的话可以通过 SQL 里直接做类型转换来处理比如CAST(amount AS DECIMAL(10,2))。另外一个写入冲突的场景是 mapping 冲突第一次写入的字段是long后续数据里同一字段变成了textES 会报mapper_parsing_exception。这种情况多半是源表字段类型被改过或者是不同表同步到了同一个索引但字段类型定义不一致。出现后不要急着删索引重建先确认源头数据类型再决定是调整 ES mapping 还是调整消费端的转换逻辑。6.4 死信处理与同步中断恢复ES 服务重启、网络闪断、mapping 冲突都可能导致消费端写入失败。一个可靠的消费链路必须考虑失败重试。我的做法是在消费端把失败批次记录下来解析失败的 binlog 事件、写入 ES 失败的原因、原始数据快照全部打到日志或者专门的死信表里。同时消费逻辑要支持重放即从某个位点开始重新消费。因为你手上有 binlog 的位点即使 Canal Server 那边数据已经 ack 删掉了只要 MySQL 的 binlog 文件还在看expire_logs_days的保留时间你随时可以从头重新消费一段时间的数据。这也是为什么第一节里我要强调 binlog 保留时间不能太短至少 7 天给错误恢复留够窗口。6.5 存量数据初始化与增量追平最后说一个很多新手容易忽略的问题Canal 一般只同步启动之后的新增变更存量数据它不管。所以上生产之前要先做一次全量初始化和增量的衔接。推荐的步骤是先创建好 ES 索引和 mapping停掉业务写入或者选择低峰期用 DataX 或 Logstash 全量导入存量数据到 ES启动 Canal 链路从全量开始时刻之后的 binlog 增量开始消费为什么强调全量导入和 Canal 启动之间的时间衔接因为如果全量导了 10 点之前的数据Canal 却在 11 点才启动那 10 点到 11 点之间的变更就丢了。实践中我会在全量导入完成后立刻读取 MySQL 当前的 binlog position然后把这个 position 写到 Canal 的配置里再启动 Canal保证从全量结束的那个点开始无缝接续。这个细节没处理好你会看到 ES 里的数据大部分对但总差一小部分花很长时间排查元凶最后发现是增量衔接断层。我自己把这条链路完整跑通后最大的体会是Canal 本身并不复杂真正的复杂度全在消费端的可靠性和两端数据比对上。版本匹配、权限、时区这类问题都有固定解法提前按本文的思路准备好基本能避免 80% 的启动期故障。如果你们团队后续的同步场景越来越多建议推动把消费端抽象成公共组件统一处理位点、重试、类型转换和监控告警。这套做好了MySQL 到 ES 的同步就真的变成一件顺其自然的事了。
返回列表