Canal实战:基于MySQL binlog的增量数据同步到Redis与Kafka完整指南

发布时间:2026/10/6 3:29:35
Canal实战:基于MySQL binlog的增量数据同步到Redis与Kafka完整指南 想把MySQL的数据实时同步到Redis、Elasticsearch或者消息队列里去又不想在业务代码里写一堆双写逻辑更不想为了同步数据去解析那些乱七八糟的binlog格式那Canal基本是绕不开的一个组件。这是阿里巴巴开源的一套基于MySQL binlog的增量订阅与消费组件原理上就是把自己伪装成MySQL的slave节点靠主从复制协议去接收binlog事件然后解析成结构化数据再推给你想要的地方。这篇文章我会完整记录一次Canal的搭建过程从MySQL的binlog准备、Canal Server的部署、实例配置到客户端怎么消费增量数据最后再把实际踩过的坑和排查方法一并整理出来希望可以给正在调研或者打算上手Canal的人提供一套能直接照着做的方案。1. Canal为什么要存在以及它的核心定位1.1 从数据同步的痛点说起业务系统发展到一定阶段数据往往不会只待在一个MySQL实例里。订单数据要同步到搜索引擎做全文检索商品数据要同步到Redis做热点缓存运营要拉实时报表风控要盯着交易流水这些场景都需要数据库一发生变化下游系统立刻感知到变化的能力。但直接去MySQL里做轮询效率太低业务也不一定扛得住频繁的扫描查询。业务代码里做双写侵入性又太强耦合度极高一旦主流程失败数据很容易出现不一致。而且很多老系统根本不可能为了同步需求去改代码。Canal的核心思路就是绕开业务系统直接从数据库的日志层面拿数据。它通过模拟MySQL主从复制里slave节点的交互协议假装自己是一个从库去连接主库主库正常产生binlogCanal接收binlog之后做解析最终再把解析好的数据变更事件交给下游消费。这个过程对业务系统是完全透明的下游系统不需要关心上游业务怎么写的只要订阅Canal的消息就行。1.2 Canal在链路里的位置从架构上看Canal处于数据源和数据消费端之间。数据源这一侧它只适配MySQL支持的版本从MySQL 5.6、5.7、8.0到MariaDB都可以前提是开启了binlog并且使用ROW模式。消费端那一侧Canal提供了多种对接方式直接走TCP协议给Java客户端消费或者对接Kafka、RocketMQ这类消息中间件也可以输出到日志文件做离线分析。用一句话来形容Canal在数据链路里的角色它是一个非常称职的搬运工。它不负责决定数据搬到哪去也不负责下游怎么加工它只保证把MySQL的binlog变更事件完整、准确、及时地搬运到消费端。这种清晰的边界划分让Canal在实际项目里非常好落地——谁消费、消费后干什么全由业务系统自己决定。1.3 选Canal而不选其他方案的理由市面上做数据库增量同步的组件不止Canal一个比如Debezium、Maxwell还有阿里内部的DTS。我选Canal主要看重三点。第一它对MySQL binlog的解析非常成熟。Canal从2014年左右开始对外开源在阿里巴巴内部经过了大量生产环境的锤炼对binlog的格式兼容、DDL解析、主从切换处理都做得比较完善。第二它的部署形态足够轻量。一个Java进程一份配置文件改一改就能跑起来不需要依赖外部存储不像Debezium那样通常还要配合Kafka Connect框架一起使用。第三它的消费方式灵活既支持TCP直连也支持消息队列很多中小团队最需要的就是这种少一点中间环节的方案。2. 搭建Canal之前先把MySQL这头的基础打牢2.1 启动增量数据的第一步打开binlogCanal能不能工作前提条件是MySQL必须开启binlog而且日志格式必须设置成ROW。这个配置属于MySQL的必选项写在my.cnf或者my.ini里。[mysqld] log-binmysql-bin binlog_formatROW binlog_row_imageFULL server-id1这里逐个解释一下这几个参数的作用。log-binmysql-bin是开启binlog并设置日志文件的前缀后面的序号由MySQL自动管理。binlog_formatROW表示让MySQL记录每一行数据是怎么变的只有ROW模式才能拿到完整的变更前后值这也是Canal解析的数据基础。binlog_row_imageFULL确保每行变更都记录全部字段的前后值如果设置成MINIMAL某些情况下只有被修改的字段有值增量数据就不完整了。server-id必须设置而且不能和复制拓扑里的其他节点冲突Canal作为伪从库连接时也会用到它。改完之后重启MySQL然后执行SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format;看到结果分别是ON和ROW就说明binlog已经正常开启了。2.2 为Canal单独创建一个同步账号Canal要伪装成从库去拉取binlog所以它需要一个数据库账号这个账号不需要业务表的任何读写权限但必须赋予复制相关的权限。我习惯单独创建账号不跟业务账号混用这样以后排查问题也清晰。CREATE USER canal% IDENTIFIED BY canal_passwd; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; FLUSH PRIVILEGES;SELECT权限是Canal在一些场景下需要回查表结构用的REPLICATION SLAVE和REPLICATION CLIENT分别是拉取binlog和查看主库状态所需的权限。注意控制好账号的网段限制生产环境最好把%换成实际部署Canal机器的IP。2.3 确认Canal要监控哪张表Canal在解析binlog时默认会订阅整个实例下所有库所有表的变化但实际生产里我们很少需要全实例的数据通常只关心某一两个业务库甚至只关心某几张表。这种过滤可以在Canal的实例配置文件里做不需要去MySQL端设置后面我会详细讲过滤规则的写法。这里要强调一个细节binlog是实例级的只要MySQL开启了ROW模式的binlog这个实例下所有库的变更都会被记录。Canal拿到这些日志后再根据自己的配置决定哪些数据需要下发、哪些可以丢掉。所以从数据安全角度考虑binlog会带来一定的磁盘空间增长开启之前要预估好增长速度尤其是做表结构变更、大批量update的时候ROW模式产生的binlog体积会明显变大。2.4 消费端基础设施的选择Canal解析出来的增量数据最终要交给下游。下游用哪套方式消费决定了我需要在搭建前准备什么基础设施。如果业务是Java技术栈下游需要的是同步接口调用可以直接用Canal自带的TCP模式客户端通过CanalConnector连接Canal ServerPush模式消费增量消息。这种模式最简单不需要额外引入消息队列。但如果下游有多个系统都要消费同一份增量数据或者想利用消息队列的重试、堆积能力那Kafka或者RocketMQ会更合适。我在实际项目中大多数情况走Kafka因为公司已有的消息链路基本都是Kafka接入成本低而且Canal对Kafka的适配做得很好支持自动创建topic、分区顺序保证这些关键能力。版本选择方面我推荐用1.1.x系列比如1.1.7。这个版本对MySQL 8.0的支持比较稳定也支持了自定义时间戳、同步进度回调等功能小版本迭代过程中的坑相对更少。3. Canal Server的部署和配置3.1 准备好运行环境Canal Server本身是一个Java应用官方提供了解压即用的发行包也提供了Docker镜像。我优先推荐用发行包部署部署过程中可以看到完整的日志输出出问题也更容易定位。运行环境需要JDK 1.8及以上我一般用OpenJDK 8内存至少给2G因为Canal在解析大事务、大批量binlog时会有比较大的内存占用。下载的时候要注意选择canal.deployer-1.1.7.tar.gz不要下成canal.adapter或者canal.adminadmin是Web管理台adapter是配合做异构同步的适配器我们这里只需要最核心的deployer。解压之后目录结构大概是这样的canal.deployer-1.1.7/ ├── bin ├── conf │ ├── canal.properties │ ├── logback.xml │ ├── spring │ │ └── file-instance.xml │ └── example │ └── instance.properties └── libconf目录下的canal.properties是Canal Server的全局配置example目录下的instance.properties是具体某个数据采集实例的配置。Canal Server可以同时运行多个instance每个instance对应一个MySQL数据源。3.2 canal.propertiesServer级别的关键配置先看canal.properties用编辑器打开后重点需要关心三类配置。第一类是服务端口和模式。默认情况下Canal会同时开启两个端口canal.port11111是TCP消费端口客户端通过这个端口连接canal.admin.port11110是管理端口。如果走Kafka模式canal.serverModeKafka这时候TCP端口就不具备实际消费意义了。第二类是注册中心相关的配置。1.1.x版本的Canal把instance的配置信息抽象成了MetaManager不管是不是集群模式都需要指定注册中心的实现类。如果只是单机部署用默认的MemoryMetaManager就行如果需要多机共享配置可以放在ZooKeeper里。第三类是消费端对接配置。走Kafka模式时需要告诉Canal Kafka的地址、topic的命名规则、分区数等。这些配置有的写在全局有的可以在instance里覆盖。下面是一个走Kafka模式的最小化配置canal.serverMode kafka canal.port 11111 canal.zkServers canal.instance.global.spring.xml classpath:spring/default-instance.xml canal.mq.servers 127.0.0.1:9092 canal.mq.producerGroup canal_group canal.mq.flatMessage truecanal.mq.flatMessage true这部分要注意它决定了消息体的编码格式。true的时候Canal会把解析出的数据转成扁平化的JSON字符串生产消费都很直观适合大多数场景false的时候会走protobuf序列化格式体积更小但消费端需要额外做反序列化除非对性能特别敏感否则我建议保持true。3.3 instance.properties采集实例的完整配置instance.properties是整个Canal配置里最需要小心对待的文件它决定了Canal从哪个MySQL实例采集数据、用什么账号登录、过滤哪些表、从什么位点开始消费。# MySQL地址 canal.instance.master.address 127.0.0.1:3306 # binlog日志位点 canal.instance.master.journal.name canal.instance.master.position canal.instance.master.timestamp # 数据库账号 canal.instance.dbUsername canal canal.instance.dbPassword canal_passwd # 过滤规则 canal.instance.filter.regex test\\.user.* canal.instance.filter.black.regex # 表结构缓存 canal.instance.parser.support.ddl truemaster.address是MySQL主库的地址和端口这里要注意Canal连接的是主库还是从库。如果只是想采集数据连接从库也可以毕竟从库同样接收binlog并以相同的格式写入本地文件。但是如果从库存在复制延迟那Canal拿到的变更就会延后。我的习惯是同步业务要求实时性高就直连主库只做备份或者实时性不敏感的再考虑从库。master.position是位点配置如果这里留空Canal启动后会从MySQL当前最新的binlog位点开始消费。但实际生产环境往往会遇到存量数据已经存在只希望采集从部署时刻之后的增量的场景这种场景留空就足够了。如果遇到Canal之前挂掉了需要从某个之前消费到的位置继续的场景就要手动指定journal.name和position这两个值可以从MySQL执行SHOW MASTER STATUS拿到。过滤规则canal.instance.filter.regex的语法是Schema名加表名中间用\\.分隔支持通配符。多张表用逗号分隔比如canal.instance.filter.regex test\\.user, order\\..*, goods\\.sku表示监听test库下的user表、order库下的所有表、goods库下的sku表。如果只想要一张表一定要写清楚完整路径否则过滤太宽会把不需要的表数据也拉进来。3.4 启动Canal Server配置完成后就可以启动部署了。启动命令在bin目录下Linux环境执行sh bin/startup.sh启动过程会创建一个后台守护进程进程日志默认写在logs/canal/canal.log里每个instance的日志写在logs/example/example.log里。启动完第一步不是急着看数据有没有推送而是先确认两件事。第一看canal.log里有没有报错正常启动会打印Canal启动完成、端口监听的记录。第二看example.log里有没有出现success to load binlog position之类的信息这说明Canal已经成功找到位点、开始伪装成从库拉取binlog了。如果启动后日志里出现连接超时、认证失败基本就是MySQL地址、账号权限或者网络不通的问题直接用客户端工具连一下MySQL验证即可。4. 客户端接入把增量数据真正用起来4.1 TCP模式下的Java客户端写法和消息结构Canal最容易上手的使用方式是TCP直连适合增量数据量不大、下游逻辑不复杂的Java应用。CanalConnector connector CanalConnectors.newSingleConnector( new InetSocketAddress(127.0.0.1, 11111), example, , ); connector.connect(); connector.subscribe(test\\\\.user.*); while (running) { Message message connector.getWithoutAck(100, 1000); long batchId message.getId(); if (batchId -1 || message.getEntries().isEmpty()) { Thread.sleep(1000); continue; } for (CanalEntry.Entry entry : message.getEntries()) { if (entry.getEntryType() CanalEntry.EntryType.ROWDATA) { CanalEntry.RowChange rowChange CanalEntry.RowChange.parseFrom(entry.getStoreValue()); System.out.printf(table:%s, event:%s%n, entry.getHeader().getTableName(), rowChange.getEventType()); } } connector.ack(batchId); } connector.disconnect();这段代码里的核心流程是connect建立连接、subscribe订阅感兴趣的表、循环调用getWithoutAck批量获取数据处理完后再调用ack确认消费。getWithoutAck的第二个参数是超时时间意思是最多阻塞多久返回一批数据这样循环里不至于空转太频繁。这里必须强调ack和rollback的用法。如果一批数据处理到一半失败了不要调用ack应该调用connector.rollback(batchId)Canal会把这一批数据重新投递。这保证了消息的至少一次语义也就是说消费端要做好幂等处理因为同一个事件在异常重试场景下可能会收到两遍。消息结构里最关键的类是CanalEntry.RowChange它包含了事件类型、表名以及变更前后每一列的值。列的数据放在RowData里beforeColumns是修改前的快照afterColumns是修改后的快照通过遍历这两组列再按列名拼装就能还原出完整的数据变更。4.2 消息格式的两种选择前面提到flatMessage设为true时Kafka里的消息体就是一个扁平的JSON字符串。这是我首推的方式因为消费端完全不需要引入Canal的客户端依赖任何语言都能消费。消息体长这个样{ data: [ { id: 102, name: 测试商品, price: 19.9 } ], database: shop, es: 1710146807000, id: 6, isDdl: false, old: [ { price: 9.9 } ], pkNames: [id], sql: , table: product, ts: 1710146809000, type: UPDATE }字段含义比较直观type是事件类型database和table指明来源data是当前数据old是变更前的旧值pkNames是主键字段列表es和ts分别表示事件发生时间和Canal处理时间。消费端拿到这条消息判断type是INSERT、UPDATE还是DELETE再结合主键去更新Redis里的缓存或者调用搜索引擎的更新接口一套简易的实时同步就通了。4.3 Kafka模式下从零开始消费走Kafka模式时Canal Server启动过程中会自动在建好的topic目录下创建分区并开始推送消息。消费端只需要按标准的Kafka消费组来订阅那个topic。Component public class CanalKafkaConsumer { KafkaListener(topics example_product, groupId canal-demo) public void onMessage(String message) { JSONObject json JSON.parseObject(message); String type json.getString(type); JSONArray data json.getJSONArray(data); if (INSERT.equals(type) || UPDATE.equals(type)) { for (int i 0; i data.size(); i) { JSONObject row data.getJSONObject(i); String key product: row.getString(id); redisTemplate.opsForValue().set(key, row.toJSONString()); } } else if (DELETE.equals(type)) { for (int i 0; i data.size(); i) { JSONObject row data.getJSONObject(i); String key product: row.getString(id); redisTemplate.delete(key); } } } }这里有几个容易忽视的点。Canal投递到Kafka时是按pkNames里的字段做分区key的所以同一行数据的变化都会进同一个分区消费端读取时天然能保证单行数据的顺序。如果有跨行的事务需求比如同一个事务里改了多张表Kafka模式下这些条目虽然带着相同的事务ID但投递到不同分区后无法保证读取顺序这种场景就需要在消费端自己根据事务ID做聚合不过大多数同步场景不会涉及这么苛刻的要求。另外记得Kafka消费组里的消费者数量最好不要超过topic的分区数否则多出来的消费者会一直处于空闲状态。Canal创建topic时默认分区数是canal.mq.partitions参数指定的值如果消费TPS要求高把这个值调大一些就可以了。4.4 通用同步抽象一套写好的消费模板看多了同步场景之后你会发现不管下游是Redis、ES还是别的存储消费逻辑都可以抽象成一个模板解析消息类型取出主键取出业务字段执行相应的更新动作。我自己在项目里习惯把消费模板做成一个统一的接口每个同步目标只实现自己的处理器。比如定义一个SyncHandler接口包含onInsert、onUpdate、onDelete三个方法每个下游系统的适配器实现这些方法。这样当需要新增一个同步目标时不需要改Canal的接入代码只需要增加一个新的Handler实现类。这种设计看起来多写了几个类但长期维护的收益非常大尤其是Canal存量表很多、下游系统更多的时候一套模板能避免大量重复代码。5. 生产环境里最容易踩的那些坑5.1 binlog模式不对导致解析报错经常遇到的情况是MySQL已经跑了一段时间binlog_format还是默认的STATEMENT或者MIXED模式Canal启动后会报类似canal parse error或者解析出的数据与预期不符。排查方法就是回MySQL去看当前binlog格式如果格式不是ROW单独改binlog_formatROW还不够需要让新配置对所有后续连接都生效并且确认Canal连接的会话用的确实是ROW。改完配置文件后重启MySQL服务再确认一遍当前值。这里还要注意即使binlog_format改成了ROW之前已经生成的binlog文件依然是旧格式Canal如果从旧的位点开始消费依然可能解析失败。所以切换格式之后最好是清空Canal的消费位点让它从当前最新的binlog开始而不是从旧位点追。5.2 位点丢失和数据重复Canal的位点管理存在内存里也支持配置持久化。如果Canal进程突然宕机重启后如果用的是默认位点管理它可能会从内存记录的最后位点继续但如果配置的是从当前最新位点消费那就有可能在宕机期间漏掉一部分数据。解决思路有两个方向。一是在Canal端开启位点定时持久化把位点信息写到文件中这样重启后能从文件恢复二是从消费端兜底因为Canal给出的是至少一次的消费语义消费端只要保证处理逻辑幂等即使偶发重复数据也不会造成大问题。我在实际生产里是两种同时做Canal端开启持久化消费端每个处理逻辑都设计成按主键幂等。5.3 大事务导致的内存压力ROW模式下一个事务里update了一百万行binlog就会有一百万行的变更记录。Canal拿到这个大事务时会先解析并暂存在内存里再批量投递到下游。如果事务特别大Canal JVM的堆内存设置不够会出现明显的Full GC甚至OOM。对这种场景常规手段是调大Canal的JVM内存JVM参数在bin/startup.sh里通过JAVA_OPTS指定。此外尽量不要在业务高峰期做超大批量的update或delete操作从源头控制binlog的增长速度。如果历史数据有大量变更要做建议拆批执行每批几千行既减轻MySQL的压力也减轻Canal的负担。5.4 DDL变更对同步链路的影响Canal默认会把DDL语句也发布到消费端比如ALTER TABLE、CREATE TABLE等。在flatMessage模式下DDL消息里isDdl字段是truedata字段为空sql字段会带着完整的DDL语句。如果消费端没有对DDL做处理按正常数据消息的逻辑去解析就会出错或者空指针。我处理DDL的原则很简单如果下游是Redis这类缓存DDL消息直接忽略掉因为缓存结构是独立的业务代码自己控制结构变更如果下游是ES这类需要定期对齐数据库表结构的系统DDL消息需要走一个告警或者手工处理的通道避免数据库表结构改了、ES mapping没跟上导致后续数据写入失败。5.5 主从切换后的位点适配MySQL做高可用切换后原来的主库变成了从库新的主库上Canal之前消费到的binlog文件名和位点可能对不上导致Canal拿到新主库上不存在的位点信息。此时最稳妥的处理是让Canal重新从新主库的最新位点开始消费代价是切换期间产生的增量数据会丢失。对增量同步链路来说切换期间的少量丢失在有些场景可以接受但有些场景不行。如果业务对完整性要求极高建议在MySQL高可用切换前先把Canal进程停掉等主从切换完成、新主库稳定后再配置新主库地址和最新位点重新启动Canal。这样丢失窗口只出现在Canal停止到新主库就绪之间的时间段。6. 一次完整的实战复盘从建库到Redis双写缓存6.1 业务背景和整体目标之前有一个电商项目商品表存在MySQL里redis里存放商品详情缓存。最初的做法是业务代码里修改商品时同步更新Redis缓存看起来简单但后来加了多个修改入口之后发现总有漏掉的地方缓存和数据库经常不一致。于是决定引入Canal做最终的兜底同步任何表数据变更最终都会通过Canal同步到Redis所有业务方不需要再关心缓存更新只管写数据库。6.2 表结构和链路配置商品表shop.product核心字段包括id、name、price、stock、status。目标把这张表的实时变化同步进Rediskey设计为product:{id}value为商品JSON。MySQL侧确认binlog开启且是ROW模式创建canal账号。Canal Server配置一个instancefilter设置成shop\\.product消费模式走Kafkatopic命名为canal-product-topicflatMessage开启。消费程序是Spring Boot应用监听Kafka消息根据type更新Redis。6.3 实际运行效果和观察启动后我在数据库做了一次update操作把某个商品价格从19.9改成29.9大概几百毫秒内Redis里的商品缓存就变成了新值。连续做了insert、update、delete三组测试Redis里的数据都跟着正确变化。确认链路通之后再观察topic的消息上报确认每个操作都对应一条消息消息里的数据都完整。这套链路上线后的收益很明显业务方删掉了几处手动更新缓存的代码减少了业务逻辑和缓存逻辑的耦合数据变化后不管是从后台系统改的、定时任务改的、还是运营手工在数据库改的只要写进MySQL缓存最终都会同步。这就是增量订阅组件的价值所在。6.4 回顾搭建过程中的决策回看这个项目有几个决策值得记录。为什么会先把消息打进Kafka而不是直接TCP给消费程序因为考虑到以后可能不止一个下游系统需要这些变更Kafka做成一个独立的数据管道未来接ES、接数仓都不用改上游。为什么flatMessage用true而不是protobuf因为消费端和CanalServer不是同一套版本也没关系拿JSON自己解析更灵活。这些决策本身没有绝对的对错关键是要符合自己的业务预期。7. 消费端的进阶技巧和设计建议7.1 幂等设计怎么做最省心Canal的消费语义决定了数据可能重复所以消费端一定要做幂等。我的做法是给每条处理逻辑都定义一个业务主键更新Redis时直接用这个主键设值天然幂等更新数据库时先查再更或者用ON DUPLICATE KEY UPDATE。总之不要假设每条消息只会来一次。7.2 压测时关注哪些指标搭建完成后最好做一轮简单的压测来评估链路能力。需要关注三个指标一是Canal Server的吞吐量单位时间内处理了多少binlog事件二是Kafka消费端Lag确认消费能力有没有成为瓶颈三是同步端到端延迟从数据库落库到下游消费完成的时延。这三个指标都能在Canal日志、Kafka监控和消费程序的日志里找到。7.3 通用同步框架的扩展路径Canal本身只做增量订阅这一件事但如果想构建一套完整的异构数据同步平台还需要补齐几个模块元数据管理、同步任务配置、监控告警、全量迁移工具。Canal社区的开源方案里常搭配Canal Adapter做异构同步比如把数据同步到ES、HBase等但我更建议按需组合简单的场景用Canal加Kafka就够了复杂场景再考虑引入Adapter。8. 写在最后个人实操经验小结Canal给我的感觉一直是一个部署简单但细节很多的组件。部署本身半小时就能跑通但真正决定链路稳不稳的全在对位点、数据格式、幂等和异常处理这几个细节的把握上。这里分享几点我在实际项目里验证过的建议MySQL的binlog格式一定提前确认是ROWCanal账号权限给到位先做小范围表测试、连通后再扩展范围消费端务必设计好幂等宁可做得保守一些消息格式优先选flatMessage运维和调试都省心测试阶段养成定期查看canal.log和instance日志的习惯很多问题是慢慢积累而不是突然爆发的。对于刚接触Canal的人我的建议是先从一个小而完整的场景入手比如只同步一张商品表到Redis把整条链路跑通、把消息格式看明白再逐步扩展到更多表和下游系统。这样即使遇到问题排查范围也是可控的。增量订阅这件事做得好能让业务系统省掉大量重复代码也让数据一致性有了一个兜底保障做得不好反而会引入新的故障点。所以每一条配置、每一个消费逻辑都值得你在上线前多想一步。