基于Canal与MySQL Binlog的实时数据同步方案设计与实践
1. Canal与MySQL Binlog实时数据同步方案设计第一次接触Canal是在2018年做电商系统重构时当时需要解决订单数据和库存数据的实时同步问题。传统的定时任务拉取方式不仅延迟高还给数据库带来不小压力。经过技术调研最终选择了基于Canal的解决方案效果出乎意料的好。1.1 为什么选择Canal在实际项目中我们遇到过几种常见的数据同步需求订单生成后需要实时更新ES搜索引擎用户信息变更需要同步到CRM系统商品价格调整需要立即刷新缓存早期我们尝试过以下几种方案定时任务扫描设置每分钟扫描变更表但总有延迟且浪费资源触发器消息队列对数据库性能影响较大维护成本高双写业务代码中同时写两个库容易产生数据不一致相比之下Canal的方案有三大优势实时性毫秒级延迟业务几乎无感知低侵入不修改业务代码对数据库压力小灵活性可以对接多种下游系统Kafka、ES、HBase等1.2 整体架构设计一个完整的实时同步系统通常包含以下组件MySQL Server │ ▼ Binlog文件 │ ▼ Canal Server伪装Slave │ ▼ Canal ClientJava应用 │ ▼ Kafka/RocketMQ可选 │ ▼ 数据消费端ES/HBase/其他DB我在金融项目中实践过的典型配置MySQL 5.7开启GTID模式Canal 1.1.5版本集群部署通过Kafka中转消息最终写入Elasticsearch和HBase双存储这种架构每天能稳定处理千万级数据变更峰值时延控制在500ms内。2. 环境搭建与配置实操2.1 MySQL环境准备要让Canal正常工作MySQL需要做好这些准备开启binlog这是最关键的步骤# 查看当前binlog状态 SHOW VARIABLES LIKE log_bin; # 在my.cnf中添加配置 [mysqld] log-binmysql-bin binlog-formatROW server_id1 expire_logs_days3创建Canal专用账号CREATE USER canal% IDENTIFIED BY Canal123; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; FLUSH PRIVILEGES;验证binlog格式SHOW GLOBAL VARIABLES LIKE binlog_format;必须确保是ROW模式其他模式会导致无法获取完整数据变更。2.2 Canal服务部署我推荐使用Docker方式部署省去环境依赖的麻烦docker pull canal/canal-server:v1.1.5 docker run -p 11111:11111 --name canal \ -e canal.instance.mysql.slaveId1234 \ -e canal.instance.filter.regex.*\\..* \ -d canal/canal-server:v1.1.5关键配置项说明instance.properties# 主库地址 canal.instance.master.address127.0.0.1:3306 # 账号信息 canal.instance.dbUsernamecanal canal.instance.dbPasswordCanal123 # 监听规则所有库所有表 canal.instance.filter.regex.*\\..* # 批量获取大小 canal.instance.mysql.batch.size10003. 客户端开发实践3.1 基础Java客户端这里分享一个经过生产验证的客户端代码模板public class CanalClient { public static void main(String[] args) { // 创建连接 CanalConnector connector CanalConnectors.newSingleConnector( new InetSocketAddress(127.0.0.1, 11111), example, , ); int batchSize 1000; try { connector.connect(); connector.subscribe(.*\\..*); while (true) { Message message connector.getWithoutAck(batchSize); long batchId message.getId(); if (batchId ! -1) { processEntries(message.getEntries()); connector.ack(batchId); // 确认消息 } Thread.sleep(1000); } } finally { connector.disconnect(); } } private static void processEntries(ListEntry entries) { for (Entry entry : entries) { if (entry.getEntryType() EntryType.TRANSACTIONBEGIN || entry.getEntryType() EntryType.TRANSACTIONEND) { continue; } RowChange rowChange RowChange.parseFrom(entry.getStoreValue()); EventType eventType rowChange.getEventType(); for (RowData rowData : rowChange.getRowDatasList()) { if (eventType EventType.DELETE) { printColumn(rowData.getBeforeColumnsList()); } else if (eventType EventType.INSERT) { printColumn(rowData.getAfterColumnsList()); } else { System.out.println(------- 更新前); printColumn(rowData.getBeforeColumnsList()); System.out.println(------- 更新后); printColumn(rowData.getAfterColumnsList()); } } } } }3.2 高级应用技巧批量处理优化// 调整获取批次大小 Message message connector.getWithoutAck(5000); // 使用并行流处理 entries.parallelStream().forEach(this::processEntry);断点续传实现// 保存position到Redis String positionKey canal:position: destination; redisTemplate.opsForValue().set(positionKey, batchId); // 重启后恢复 Long lastBatchId redisTemplate.opsForValue().get(positionKey); if (lastBatchId ! null) { connector.rollback(lastBatchId); }监控告警集成// 使用Micrometer指标 Metrics.counter(canal.event.count).increment(entries.size());4. 生产环境最佳实践4.1 高可用方案设计在电商大促场景下我们是这样保证稳定性的双Canal服务架构MySQL Master ├── Canal Server A机房A └── Canal Server B机房B ├── Kafka Cluster └── 多个消费组关键配置两个Canal使用不同的slaveId开启GTID模式避免位置冲突Kafka设置3副本保证消息不丢失4.2 常见问题排查延迟高问题检查SHOW SLAVE STATUS中的Seconds_Behind_Master调整canal.instance.mysql.batch.size参数增加Canal服务器CPU资源内存溢出处理# 调整JVM参数 canal.manager.jvm.opts-Xmx4096m -Xms4096m网络闪断恢复// 客户端增加重试逻辑 for (int i 0; i 3; i) { try { connector.connect(); break; } catch (Exception e) { Thread.sleep(5000); } }4.3 性能优化建议根据压测结果总结的优化点参数调优canal.instance.mysql.batch.size2000 canal.instance.mysql.batch.timeout500 canal.instance.transaction.size1024资源分配建议每1000TPS需要分配1核CPU每条消息按1KB计算内存需求网络带宽 TPS * 消息平均大小监控指标延迟时间canal.delay.time处理速率canal.event.rate错误次数canal.error.count