
做后端这几年被缓存和搜索引擎的数据同步折腾过不少次。典型的场景就是MySQL 是唯一业务数据源订单、用户、商品全都存在里面用户请求一多MySQL 扛不住于是 Redis 里披了各种热点缓存同时运营要做搜索、筛选ES 那边还挂着一份供检索用的索引副本。每次 MySQL 里的数据一改Redis 缓存和 ES 索引都得跟着变。变慢了用户看到的是旧数据变漏了排查半天都找不到根因。Canal 解决的就是这个痛点——它监听 MySQL 的 binlog 日志把数据库的每一次变更实时解析出来再推给下游业务由我们决定去更新 Redis 还是刷新 ES。这套方案对业务代码零侵入不用在每条 SQL 后面手动补同步逻辑这也是我把整条链路稳定跑起来之后最想推荐给同行用的原因。1. 数据同步需求拆解与方案选型1.1 业务场景里的数据一致性问题我用一个实际订单系统举例。订单表 order_info 存数据库App 端首页展示我的订单走 Redis 加速后台运营要按订单号、商品名搜索查的是 ES。某个用户支付成功订单状态从 0待支付变成 1已支付。如果只改了 MySQLRedis 里还是旧状态App 刷新后也不对ES 里还是未支付导致运营搜不到这批订单这直接影响用户信任和营收。数据不一致有三个来源缓存写失败、消息通知丢失、同步任务延迟。如果采用最简单的事务后手动更新缓存代码里到处散落同步逻辑系统一复杂就漏用定时任务全量刷新缓存和索引短时间数据量大时MySQL 压力大实时性也不够。这类问题做架构的同行应该都有体会。1.2 常见同步方案对比我把几种常见做法放在一起对比过各有各的适配场景方案实现成本实时性侵入性风险点业务代码双写低高极高业务代码中夹杂缓存更新逻辑跨系统事务难保证漏写一处就出脏数据定时任务全量刷新低低分钟级起步低但定时任务本身要维护数据量大时 DB 和 ES 压力大刷写窗口内有短暂不一致应用发 MQ 通知下游中高中依赖事务消息或本地消息表事务提交前后消息顺序难保证MQ 故障时消息丢失Canal 监听 binlog较高高秒级甚至毫秒级低业务完全无感知依赖 binlog 开启运维能力要求稍高Canal 最大的优势是“从数据根源解决问题”。数据库的 binlog 天然记录每一次增量变更Canal 不侵入业务、不占用业务线程只是把自己伪装成 MySQL 的 slave 去拉取日志。同时它天然捆绑了事务的先后顺序在解析端能拿到完整的前后镜像这对于构建缓存和索引来说信息量足够。我在实际项目里选择 Canal还有一个隐性原因业务系统历史悠久不可能为了数据同步把每条写 SQL 都改一遍。引入 Canal 后业务侧代码零改动只加一个独立的数据管道服务风险和成本都可控。2. 原理基础binlog 与 Canal 的工作机制2.1 binlog 日志本质与开启方式binlog 是 MySQL 的二进制日志它记录的是数据库所有变更操作的“流水账”比如哪张表插了一行、哪一列更新成什么值都有完整记录。MySQL 主从复制就是在一台备用 MySQL 上重放主库的 binlog最终实现数据一致。使用 Canal 的前提是 MySQL 必须开启 binlog并且要求binlog_formatROW。配置写在 MySQL 的 my.cnf 或 my.ini 中[mysqld] server-id1 log-binmysql-bin binlog_formatROW binlog_row_imageFULL expire_logs_days7MySQL 8.0 之后推荐用binlog_expire_logs_seconds控制清理时间比如 604800 秒表示7天。修改配置后重启 MySQL用下面的 SQL 验证SHOW VARIABLES LIKE log_bin%; SHOW VARIABLES LIKE binlog_format%; SHOW MASTER STATUS;SHOW MASTER STATUS能查询到当前 binlog 文件名和 Position这个位置信息是 Canal 启动订阅的起点后面做全量增量衔接时会用到。2.2 为什么 Canal 强制要求 ROW 格式binlog 有三种格式格式记录内容使用场景STATEMENT记录执行的 SQL 语句日志量小但非确定性函数会导致主从不一致ROW记录每一行的实际数据变更最详细能拿到每行数据的前后镜像Canal 依赖它MIXED混合模式自动选择以上两种兼容性好但无法保证 Canal 拿到完整行数据Canal 解析 binlog 时需要从日志中还原出完整的变更数据。STATEMENT 格式只存 SQLCanal 无法推断每一行改动后的具体值MIXED 格式则会在部分场景退化成 STATEMENT。所以官方建议统一binlog_formatROW同时把binlog_row_image设为 FULL确保 update/delete 时 MySQL 日志里包含所有列的旧值镜像Canal 解析出的数据才是完整的。注意如果只把binlog_format调成 ROW但binlog_row_image保持默认的 FULL 就是 OK 的。在生产环境改 binlog 配置务必评估磁盘空间ROW 格式日志量一般会比 STATEMENT 大 3~5 倍。2.3 Canal 模拟 MySQL slave 协议的完整链路Canal 的工作原理可以这样理解MySQL 官方主从复制是备用 MySQL 实例向主库发送 dump 协议拉取 binlog 并在本地重放。Canal 做的事情几乎一样只不过它不把日志重放到 MySQL而是解析成标准数据结构发到自己的服务端再等客户端消费。整个链路分四步Canal Server 模拟 slave 向 MySQL master 发送 dump 请求。MySQL master 将 binlog 增量推送给 Canal。Canal 解析 binlog 二进制协议提取库名、表名、事件类型INSERT/UPDATE/DELETE/DDL、变更前镜像和变更后镜像。Canal Server 将解析结果包装成 Message暴露给外部客户端客户端消费并执行自定义逻辑比如写 Redis、写 ES。这个机制完全绕开了业务应用Canal 不依赖任何业务框架一旦装好业务系统依然走原来的 JDBC 访问 MySQL只是多了一个“影子从库”在暗暗订阅日志。3. 环境准备从零搭建可运行的同步链路3.1 MySQL 账号授权与连接准备Canal 需要连接 MySQL 拉取 binlog所以要建一个专用账号并赋予复制权限。注意不能用主业务的账号权限隔离更安全CREATE USER canal% IDENTIFIED BY canal_passwd; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; FLUSH PRIVILEGES;SELECT权限用于解析端读取表结构信息比如字段映射REPLICATION SLAVE是拉取 binlog 的核心权限REPLICATION CLIENT用于查询 master 状态。如果 Canal 的权限不足启动日志会直接报Access denied排查起来很直接。3.2 Canal Server 部署与核心配置我以 Canal 1.1.x 为例。下载解压后目录结构大致是conf/canal.properties服务端全局配置定义实例列表、端口号、zk、内存等。conf/example/instance.properties单个数据同步实例的配置决定订阅哪个 MySQL、哪些库表。bin/startup.sh启动脚本。canal.properties中最重要的几项canal.port11111 canal.destinationsexample canal.instance.global.modemanager canal.instance.global.manager.address127.0.0.1:8080 canal.instance.tsdb.enabletrue如果是单机测试可以不启用 admin直接指定 instance。canal.port11111是客户端消费消息的 TCP 端口。再打开conf/example/instance.properties配置订阅源库canal.instance.master.address127.0.0.1:3306 canal.instance.dbUsernamecanal canal.instance.dbPasswordcanal_passwd canal.instance.connectionCharsetUTF-8 canal.instance.filter.regextest\\..*canal.instance.filter.regex用的是 Java 正则表达式并且需要注意字符串转义。test\\..*表示订阅 test 库下的全部表如果要订阅所有库所有表写作.*\\..*。这个正则不配置的话默认订阅所有实际生产建议按需订阅减少无效日志解析。启动后看日志bin/startup.sh tail -f logs/example/example.log日志中会出现load canal process successfully之类的字样说明同步管道已经建立。用客户端连上去消费就行。3.3 验证 Canal 服务已成功接收 binlog很多同学部署完不知道 Canal 是否真的在运行我建议先用 SQL 手动改一条数据再观察 Canal 日志有没有出现Dump相关记录或者直接用 Java 客户端测试拉取消息。Canal 没有自带图形管理界面最简单的验证就是消费端能拉到Entry数据。实操心得Canal 的日志非常详细如果instance.properties的账号密码或 filter 正则写错它不会立刻退出而是反复报错重连。先看tail -f logs/example/example.log再针对报错处理通常比盲改配置高效得多。3.4 Redis 与 ES 环境准备Redis 和 ES 的部署方式我不过多展开但强调几个关键点。Redis 建议部署双副本以上用 Docker 一把拉起主从也够开发环境用。ES 8.x 默认开启了安全认证开发环境里可以把xpack.security.enabled调成 false 降低门槛生产则必须保持开启并且用专用账号访问。准备一个可视化工具也会节省很多时间Redis 可以用 Another Redis Desktop ManagerES 用 Kibana 的 Dev Tools直接写 REST 请求验证索引数据比在代码里层层排查效率高得多。4. 核心代码Java 客户端消费 Canal 消息4.1 引入 Canal 客户端依赖Canal 官方提供了 Java 客户端 SDK项目里加 Maven 依赖即可dependency groupIdcom.alibaba.otter/groupId artifactIdcanal.client/artifactId version1.1.7/version /dependency依赖中会传递 protobuf 相关组件Canal 服务端与客户端的消息协议就是基于 protobuf 封装。不推荐自己用 TCP 裸连协议细节较多直接用官方客户端最稳妥。4.2 建立连接与订阅过滤规则连接 Canal Server 的代码非常简洁CanalConnector connector CanalConnectors.newSingleConnector( new InetSocketAddress(127.0.0.1, 11111), example, , ); connector.connect(); connector.subscribe(test\\\\.user); connector.rollback();subscribe方法第二个参数是订阅过滤表达式这里的test\\\\.user表示订阅 test 库的 user 表。注意 Java 字符串中的转义正则里的\.在 Java 字符串里要写成\\.所以最终看到四个反斜杠。如果订阅规则不匹配连接是成功的但永远拉不到任何消息。4.3 解析 RowChange 数据模型Canal 消息的数据模型分为三层Message一批批量消息容器内部含多个Entry。Entry一次事务或单行变更的载体头部记录 schema 名、表名、事件类型。RowChange由Entry.getStoreValue()反序列化得到内含多行变更数据。解析示例Message message connector.getWithoutAck(batchSize); long batchId message.getId(); int size message.getEntries().size(); if (batchId -1 || size 0) { Thread.sleep(1000); } else { for (CanalEntry.Entry entry : message.getEntries()) { if (entry.getEntryType() ! CanalEntry.EntryType.ROWDATA) { continue; } String schemaName entry.getHeader().getSchemaName(); String tableName entry.getHeader().getTableName(); CanalEntry.RowChange rowChange CanalEntry.RowChange.parseFrom(entry.getStoreValue()); CanalEntry.EventType eventType rowChange.getEventType(); for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { ListCanalEntry.Column beforeColumns rowData.getBeforeColumnsList(); ListCanalEntry.Column afterColumns rowData.getAfterColumnsList(); // 业务处理 } } } connector.ack(batchId);拿到beforeColumns和afterColumns后可以封装一个辅助方法把 Column 列表转成 Mapprivate static MapString, String toColumnMap(ListCanalEntry.Column columns) { MapString, String map new HashMap(); for (CanalEntry.Column column : columns) { map.put(column.getName(), column.getValue()); } return map; }事件类型和镜像的对应关系要记牢INSERT 只有 after 有值UPDATE 同时有 before 和 afterDELETE 只有 before 有值。这在同步到 Redis 和 ES 时非常关键因为 delete 事件要用 before 主键去删除下游数据。4.4 ack 与 rollback 的事务语义Canal 消费消息和 MQ 类似必须显式确认。getWithoutAck获取一批消息处理完成后调用connector.ack(batchId)确认消费如果处理过程中抛异常调用connector.rollback(batchId)Canal 把这一批消息重新投递。在实际项目中很多人会忘记在异常分支调用rollback导致消息既没被消费也没被重发Canal 端会一直等待确认业务处理陷入僵局。我习惯在 finally 之外做判断确认逻辑尽量早把真正的幂等逻辑放在消息处理内避免重复消费影响业务。重要提示拿到batchId后第一步就要处理数据不要先睡再处理。Canal 按事务顺序推送消息如果消费端处理太慢服务端会积压大量 pending 消息。4.5 消费分发器的设计思路解析出 RowData 后不要每张表都写一个 if/else 分支建议做一个统一的分发器public void dispatch(CanalEntry.RowChange rowChange, String schemaName, String tableName) { SyncTarget target routeTable.get(schemaName . tableName); if (target null) { return; } target.sync(rowChange); }路由表可以这样组织一个 Map 维护库名.表名对应的同步策略。比如test.user表同步到 Redis 和 ESlog.order_operate_log表只同步到 ES 不做缓存更新。这样后续新增表只需配置一行不用改核心代码。5. 同步到 Redis 的落地策略5.1 缓存更新的核心思路同步到 Redis 的目的是保证用户读到的缓存与 MySQL 一致。业内最常用的是 Cache Aside 模式读请求先查 Redis未命中再查 MySQL并把结果回填写请求则直接操作 MySQL然后删除或更新缓存。传统 Cache Aside 有个明显漏洞写操作和缓存删除如果分开执行中间有窗口期会让缓存变脏。比如线程 A 写 MySQL 成功但删除 Redis 失败线程 B 立刻读到了旧缓存。Canal 天然适合解决这个问题因为 binlog 推送意味着 MySQL 已提交拿到变更事件后主动删除或重建缓存正好补上那个时间窗口。5.2 Redis 同步代码示例我常用的策略是INSERT 和 UPDATE 直接把变更后的行数据写入 RedisDELETE 删除对应缓存。这样绝大多数请求能命中最新缓存也省去了回源数据库的压力。if (eventType CanalEntry.EventType.INSERT || eventType CanalEntry.EventType.UPDATE) { MapString, String afterMap toColumnMap(rowData.getAfterColumnsList()); String key user: afterMap.get(id); String json objectMapper.writeValueAsString(afterMap); stringRedisTemplate.opsForValue().set(key, json, 30, TimeUnit.DAYS); } else if (eventType CanalEntry.EventType.DELETE) { MapString, String beforeMap toColumnMap(rowData.getBeforeColumnsList()); String key user: beforeMap.get(id); stringRedisTemplate.delete(key); }这里的 Redis key 用了user:id的结构value 存 JSON 字符串既简单又便于不同语言客户端读取。用 StringRedisTemplate 而不是 RedisTemplate避开了 jdk 序列化带来的类版本兼容问题生产环境我强烈推荐这种方案。5.3 序列化与缓存治理经验关于 Redis 序列化我踩过不少坑。RedisTemplate 默认使用 JDK 序列化value 会带 Java 类型描述信息跨端读取很痛苦。后来统一用StringRedisTemplatevalue 固定为 JSON简单直接。JSON 序列化要注意日期字段格式MySQL 的 datetime 和 Java 的 LocalDateTime 需要统一格式建议在 DTO 上使用JsonFormat明确指定。缓存治理也没有唯一标准我的经验是热点数据尽量设置一个合理过期时间避免永远不过期删除缓存时按主键删除即可不需要扫表。Canal 同步端如果批量删除建议用 Redis Pipeline 减少网络往返。实操心得更新缓存时不要直接用 MySQL 原始字段名作为 JSON 的 key 映射。比如数据库字段user_name在 Redis 和 ES 中很可能要用userName这种驼峰命名。同步链路中加一个映射层比在下游到处兼容强得多。6. 增量同步到 ES 的实现细节6.1 索引设计与 ES Java 客户端选择ES 同步的目标是让搜索索引与 MySQL 数据一致。首先是索引设计三要素id 必须和 MySQL 主键一致保证幂等字段类型要提前建好 mapping尤其是日期和 keyword 字段分片数按数据规模估算。ES 官方推荐的 Java 客户端在 7.x 时代常用RestHighLevelClient8.x 之后官方主推co.elastic.clients:elasticsearch-java也就是 Elasticsearch Java API Client。如果项目是 2026 年新起的直接使用官方新版客户端代码风格更现代dependency groupIdco.elastic.clients/groupId artifactIdelasticsearch-java/artifactId version8.13.4/version /dependencyAPI Client 的ElasticsearchClient支持链式构建请求比旧版RestHighLevelClient更简洁。6.2 ES 增量同步代码实现拿到 Canal 解析出的 RowData 后三类事件的同步逻辑分别是INSERT向 ES 索引创建文档id 使用 MySQL 主键。UPDATE用 update 请求做局部更新只传变更字段。DELETE根据主键删除文档。示例代码if (eventType CanalEntry.EventType.INSERT) { MapString, Object doc toDocMap(toColumnMap(rowData.getAfterColumnsList())); esClient.index(i - i.index(user_index) .id(String.valueOf(doc.get(id))) .document(doc)); } else if (eventType CanalEntry.EventType.UPDATE) { MapString, Object doc toDocMap(toColumnMap(rowData.getAfterColumnsList())); esClient.update(u - u.index(user_index) .id(String.valueOf(doc.get(id))) .doc(doc), Map.class); } else if (eventType CanalEntry.EventType.DELETE) { MapString, String beforeMap toColumnMap(rowData.getBeforeColumnsList()); esClient.delete(d - d.index(user_index) .id(beforeMap.get(id))); }toDocMap做的事情是把数据库字段名转成 ES 约定字段名并过滤掉不需要同步到 ES 的列。比如日志表里有content大字段没必要进索引过滤后能显著减少 ES 存储压力和网络带宽。6.3 全量与增量衔接的处理方案Canal 只解决增量如果 ES 索引刚创建必须全量导入历史数据。最稳妥的姿势分三步记录当前 MySQL master 位点SHOW MASTER STATUS记录 binlog 文件和 Position。使用 SELECT 分页查询历史数据批量写入 ES。启动 Canal并把 instance 的初始位点设置成第一步记录的位点继续订阅增量。这样全量导入期间新产生的增量数据不会丢。如果直接启动 Canal 再全量导入中间会产生重复数据ES 按 id 覆盖写入还能容忍但 Redis 缓存可能出现旧值必须按最新变更覆盖。7. 参数调优与稳定性保障7.1 消费连接与批量参数Canal 客户端的批量拉取参数直接影响吞吐。我的默认配置是batchSize1000每批最多处理 1000 行变更。太小比如 1网络往返太多太大比如 10000内存波动明显一旦处理异常滚回一整个大批次重放代价高。连接参数里soTimeout也很重要。如果消费端长时间不拉消息连接可能被服务端断开需要定期发送心跳。官方客户端内置了心跳机制但要注意消费逻辑不能阻塞太久否则客户端会认为连接异常并重连。7.2 消费端多线程与批量优化单线程消费的性能瓶颈在网络 IO 和下游写入延迟。Redis 写入很快但 ES 写入批量优化空间更大。我建议Canal 消息拉取线程只负责解析和分发不要直接写 ES。ES 写入端用队列收集变更文档满 500 条或间隔 1 秒后使用 Bulk API 批量发送。Redis 更新可以做 Pipeline 批量操作。ES Bulk 批量写入几乎能把吞吐提高数倍。实测下来单条 index 请求在非压测环境下约 5ms 一次Bulk 500 条一批时单条成本降到约 0.5ms差别是数量级的。7.3 Canal 高可用与位点恢复生产环境不能单点部署 Canal。官方支持 Canal Server 多实例配合 Zookeeper 做 HA同一个 destination 的多个实例会竞争成为 master只有 master 提供服务standby 实时等待。位点恢复是另一个重点。Canal 默认把消费位点记录在本地文件里如果 Canal Server 所在机器磁盘损坏位点丢失会引发重复或丢失消费。建议打开canal.instance.tsdb.enabletrue利用 MySQL 存储时间戳表结构支持位点回溯。这个配置在部署时直接加在instance.properties中。8. 常见问题与排查实录8.1 连接 MySQL 失败与权限报错现象一Canal 日志反复出现connect to mysql failed。排查顺序先确认instance.properties的canal.instance.master.address是否能被 Canal 机器访问再用 canal 账号手动在命令行登录 MySQL确认密码无误且权限到位。现象二日志出现Access denied for user canal。这是账号权限缺失执行GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; FLUSH PRIVILEGES;需要注意的是如果 Canal 版本较新账号插件可能要求mysql_native_passwordMySQL 8.x 默认caching_sha2_password可能导致连接失败这时可以为 canal 账号指定ALTER USER canal% IDENTIFIED WITH mysql_native_password BY canal_passwd;8.2 binlog 文件被清理导致 Canal 启动失败热词里有不少人在搜“binlog日志可以删除吗”这个问题在 Canal 场景下特别重要。如果 MySQL 的 binlog 过期清理规则太激进Canal 记录的位点可能已经指向一个不存在的日志文件启动时报错Could not find first log file name in binary log index file这种问题分两种情况处理数据允许从当前开始同步直接把 Canal 的位点清空让它从当前 master 位点开始要求不丢数据就得在 binlog 还没被清理前提前调整 MySQL 的清理配置或者定期把位点落到安全区域。我建议 binlog 保留期至少 3 天预留足够的缓冲窗口。8.3 解析时报错或字段值为空解析时最容易遇到NullPointerException或字段为空的情况。原因通常是binlog_row_image不是 FULL导致 update/delete 时beforeColumns里的列值缺失。解决办法是把binlog_row_image调整为 FULL 并重启 MySQL。还有一类情况是 DDL 事件混在变更数据里。Canal 的RowChange里有isDdl字段遇到 DDL 事件时RowData列表为空解析前必须过滤。我在代码里统一判断rowChange.getIsDdl()为 true 直接跳过。8.4 Redis 和 ES 出现脏数据脏数据多数来自重复消费或全量增量交叉。Canal 在 ack 之前的异常 rollback 会让同一批消息重新投递下游操作必须幂等。Redis 写入本身幂等DELETE 操作要判断 key 是否存在避免误删新写入的数据。ES 的 update 如果链路串行也存在旧值覆盖新值的可能建议把version字段或主键 id 带入 ES 文档用doc_as_upsert保证最终一致。8.5 性能瓶颈定位如果同步速度跟不上写入速度先看 Canal Server 日志有没有积压。再分析消费端最慢的环节通常不是解析而是下游写入。一个快速定位方法在消费端打点统计每个事件的处理耗时。我遇到过 Redis 写入慢的情况原因是用 Jedis 直连且每条都建连接换成连接池后耗时降了一个数量级。ES 写入慢则优先开启 Bulk而不是单条并发。8.6 实操经验速查表问题可能原因解决参考连不上 MySQL账号密码、网络、加密插件验权限、改 native_password拉不到消息filter 正则写错、订阅表不存在检查 subscribe 表达式解析字段为空binlog_row_image 不是 FULL修改配置并重启 MySQLConsumer 不 ack异常被吞掉或忘记 rollback异常分支调用 rollbackES 数据重复全量增量交叉或改两次按主键幂等索引吞吐低单条写入 Redis/ES改批量与连接池Canal 踩坑的通用排查步骤先看 Canal Server 日志再看消费端日志最后查目标端 Redis/ES 的实际数据。日志不会骗人只是很多人第一步就看错地方导致排查绕圈子。我个人在实际操作中的体会是这套链路最大的难点不在于 Canal 本身的配置而在于下游目标端的处理模式。Redis 缓存要理解 Cache Aside 的语义ES 索引要提前设计 mapping 和幂等策略消费端的并发模型要兼顾顺序性。踩过几次坑之后我现在的通用方案是解析层保持轻量统一分发缓存更新和 ES 更新之间用队列解耦所有写操作都以主键作为幂等键。这样无论 binlog 事件怎么重放最终数据都能收敛到一致状态。最后再分享一个小技巧把 Canal 的消费位点和业务表主键绑定在目标端记录每次同步的 binlog 位置排查数据不一致时能直接定位到哪一条变更丢了多少条比满屏日志找问题快得多。