尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Canal实战:MySQL增量数据同步与实时处理部署指南

Canal实战:MySQL增量数据同步与实时处理部署指南 1. 项目概述为什么我们需要Canal如果你正在处理一个涉及MySQL数据库的项目并且遇到了数据同步、实时分析或者缓存更新的需求那么“Canal”这个名字你大概率不会陌生。简单来说Canal是一个基于MySQL数据库增量日志解析提供增量数据订阅和消费的中间件。它的核心工作原理是伪装成MySQL的Slave通过读取主库的binlog日志将数据变更事件增、删、改实时地推送出去。这听起来可能有点抽象我举个例子你的电商平台主库每产生一笔新订单Canal就能在毫秒级内捕获到这个事件并立刻通知你的缓存服务去更新商品库存或者通知你的搜索引擎去重建索引甚至同步到另一个数据分析库去做实时报表。整个过程对主库几乎零压力实现了业务解耦和数据的实时流动。我最初接触Canal是因为一个老项目的缓存一致性难题。当时我们用的是“先更新数据库再删除缓存”的策略但在高并发下缓存删除失败或延迟会导致脏数据。引入Canal后我们让Canal监听数据库的变更由它来触发缓存的删除或更新彻底将缓存逻辑与业务代码解耦系统的稳定性和可维护性得到了质的提升。所以无论你是想构建一个实时数仓、实现多活数据中心的数据同步还是简单地想做一个可靠的缓存更新机制Canal都是一个非常值得投入学习和部署的工具。接下来我就以一个十年老兵的视角带你从零开始手把手完成Canal的安装与部署并分享那些官方文档里不会写的“坑”和技巧。2. 环境准备与前置条件核查在动手安装之前磨刀不误砍柴工把环境理清楚能避免后面80%的莫名其妙的问题。Canal的部署不是孤立的它严重依赖于MySQL数据库的配置。2.1 MySQL数据库配置重中之重Canal的本质是一个MySQL Slave所以你的MySQL必须开启Binlog并且格式必须是ROW模式。这是Canal工作的基石。首先登录你的MySQL服务器检查当前的binlog配置SHOW VARIABLES LIKE ‘log_bin‘; SHOW VARIABLES LIKE ‘binlog_format‘;如果log_bin的值是OFF说明没有开启binlog_format如果不是ROWCanal将无法正确解析数据变更的细节。你需要修改MySQL的配置文件通常是my.cnf或my.ini在[mysqld]段落下增加或修改以下配置[mysqld] # 开启binlog并指定basename这里会生成mysql-bin.000001这样的文件 log-binmysql-bin # 设置binlog格式为ROW这是Canal必需的 binlog-formatROW # 为当前服务器设置一个唯一的server-id建议用IP段或特定编号不能与其他Slave重复 server-id1 # 指定需要同步的数据库如果不指定则同步所有库。建议根据业务明确指定减少不必要的流量。 # binlog-do-dbyour_database_name修改完成后重启MySQL服务使配置生效。注意在生产环境中server-id必须全局唯一尤其是在主从复制架构中。如果是从库其server-id不能与主库或其他从库冲突。接下来你需要为Canal创建一个专门的数据库账号并授予复制权限。这个账号用于Canal连接MySQL并读取binlog。CREATE USER ‘canal‘‘%‘ IDENTIFIED BY ‘canal‘; -- 创建用户密码设为‘canal‘生产环境请使用强密码 GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO ‘canal‘‘%‘; -- 授予必要权限 FLUSH PRIVILEGES; -- 刷新权限这里解释一下这几个权限SELECT需要读取information_schema库中的元数据信息。REPLICATION SLAVE这是核心权限允许Canal以Slave身份连接并读取binlog。REPLICATION CLIENT允许使用SHOW MASTER STATUS等命令。账号创建好后务必测试一下连接确保从你计划部署Canal的机器上能用这个账号成功登录目标MySQL。2.2 Java环境准备Canal服务端和客户端Adapter Admin等都是Java应用所以需要JDK 1.8或以上版本。建议使用Oracle JDK 8或者OpenJDK 8。在服务器上执行java -version检查版本。如果没有安装可以通过系统包管理器安装如yum install java-1.8.0-openjdk-devel或者从官网下载tar包手动配置JAVA_HOME环境变量。2.3 服务器资源考量对于Canal Server本身资源消耗并不高。主要压力在于网络I/O持续读取binlog并推送给客户端。如果同步的表多、变更频繁网络流量会比较大。磁盘I/OCanal会将解析位点消费进度持久化到本地文件默认是meta.dat在高并发消费时频繁写入可能会成为瓶颈建议使用SSD磁盘。内存主要存放解析线程和事件队列。对于一般业务分配1-2GB的堆内存通过JVM参数-Xms和-Xmx设置通常足够。你可以根据监控到的GC情况和实际吞吐量进行调整。3. Canal Server的安装与部署详解Canal的部署主要分为两部分Server和Client。Server负责连接MySQL、解析binlog、管理客户端连接Client如Canal Adapter, Canal Connector则负责消费数据并写入到目标端。我们先搞定Server。3.1 获取与解压发布包Canal的官方发布在GitHub上。你可以直接下载最新的稳定版release包。这里以canal.deployer-1.1.7.tar.gz为例。# 假设我们在 /opt 目录下操作 cd /opt # 下载请替换为最新的下载链接 wget https://github.com/alibaba/canal/releases/download/canal-1.1.7/canal.deployer-1.1.7.tar.gz # 解压 tar -zxvf canal.deployer-1.1.7.tar.gz # 解压后会得到一个 canal 目录可以重命名一下 mv canal.deployer-1.1.7 canal-server cd canal-server解压后的目录结构如下canal-server/ ├── bin/ # 启动停止脚本 ├── conf/ # 配置文件目录 │ ├── canal.properties # Server全局配置 │ └── example/ # 一个实例Instance的配置目录默认叫‘example‘ │ ├── instance.properties # 该实例的详细配置连接哪个MySQL同步哪些表等 │ └── ... ├── lib/ # 依赖的jar包 ├── logs/ # 日志目录 └── plugin/ # 插件目录如解析器、Sink等3.2 核心配置文件解析与修改配置文件是Canal的灵魂理解每一个关键配置项是稳定运行的前提。首先修改全局配置文件conf/canal.properties这个文件控制Canal Server的整体行为我们关注几个关键点# canal server的工作模式tcp表示以网络端口方式提供数据kafka/rocketmq表示投递到消息队列。我们先用最简单的tcp。 canal.serverMode tcp # canal server的监听端口客户端比如我们自己的Java程序会连接这个端口来获取数据。 canal.port 11111 # 实例列表多个实例用逗号分隔。默认只有一个‘example‘对应conf/example目录。 canal.destinations example # 实例的配置自动扫描周期单位毫秒。修改instance配置后Canal会自动热加载。 canal.auto.scan true canal.auto.scan.interval 5000 # 全局的持久化机制有memory内存重启丢失、file本地文件、zookeeper集群模式。单机部署用file就行。 canal.instance.global.mode memory canal.instance.global.lazy false canal.instance.global.manager.address ${canal.conf.dir} canal.instance.global.spring.xml classpath:spring/memory-instance.xml对于单机学习或小规模使用以上默认配置基本无需改动。如果你计划部署多个实例例如连接多个不同的MySQL库就在canal.destinations里增加名字并在conf目录下创建对应的配置文件夹。然后修改实例配置文件conf/example/instance.properties这个文件定义了example这个实例具体要干什么是配置的核心。# 需要连接的MySQL地址和端口 canal.instance.master.address 127.0.0.1:3306 # 前面创建的Canal账号和密码 canal.instance.dbUsername canal canal.instance.dbPassword canal # 字符集必须和MySQL保持一致否则中文乱码 canal.instance.connectionCharset UTF-8 # 指定要订阅的数据库表这里是正则表达式。‘.*\\..*‘表示所有库的所有表。 canal.instance.filter.regex .*\\..* # 黑名单指定忽略哪些表。一般先全量再用黑名单排除系统表或不关心的表。 canal.instance.filter.black.regex mysql\\.slave_.*canal.instance.filter.regex这是最重要的过滤配置。格式为库名.表名。例如test.user同步test库的user表。test\\..*同步test库的所有表。.*\\..*同步所有库的所有表默认但生产环境慎用。canal.instance.filter.black.regex用法同上匹配到的表将被忽略。实操心得生产环境强烈建议白名单方式配置即只同步你需要的业务库表。使用全量同步.*\\..*然后黑名单排除一旦有新增的业务表忘记加黑名单Canal就会去同步可能带来不必要的网络流量和存储开销甚至引发权限问题。可以在filter.regex里明确写your_db\\.(table1|table2|table3)。3.3 启动、停止与状态检查配置完成后就可以启动Canal Server了。# 进入bin目录 cd /opt/canal-server/bin # 启动 sh startup.sh # 查看启动日志检查是否有错误 tail -f ../logs/canal/canal.log # 查看实例日志 tail -f ../logs/example/example.log正常的启动日志中你会看到Canal成功连接到MySQLcom.alibaba.otter.canal.parse.inbound.mysql.tsdb.MemoryTableMeta相关日志并开始dump binlog。停止服务使用sh stop.sh。如何检查Canal是否在正常工作查看进程ps -ef | grep canal。查看端口netstat -tlnp | grep 11111看11111端口是否在监听。查看日志重点关注example.log没有持续的报错如连接断开、解析异常通常就是正常的。简单的客户端测试可以用Canal自带的简单客户端工具canal.example进行快速测试需要单独下载client工程并运行但更常见的验证方式是直接进入下一步——部署一个客户端来消费数据。4. 客户端连接与数据消费实战Server跑起来了但它只是个“搬运工”。数据需要被消费掉才有价值。这里我介绍两种最常用的客户端模式使用原生Java客户端进行自定义处理以及使用Canal Adapter进行开箱即用的数据同步。4.1 使用原生Java客户端Canal Connector这种方式最灵活你可以在自己的Java程序里获取到每一行数据的变更事件然后做任何你想做的事情更新缓存、发送消息、记录审计日志等等。首先在你的Java项目中引入Canal Client的依赖以Maven为例dependency groupIdcom.alibaba.otter/groupId artifactIdcanal.client/artifactId version1.1.7/version !-- 版本与Server保持一致 -- /dependency下面是一个最简化的客户端示例代码import com.alibaba.otter.canal.client.CanalConnector; import com.alibaba.otter.canal.client.CanalConnectors; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.protocol.CanalEntry.*; import java.net.InetSocketAddress; import java.util.List; public class SimpleCanalClient { public static void main(String[] args) { // 1. 创建连接器连接到Canal Server的11111端口指定实例名‘example‘ CanalConnector connector CanalConnectors.newSingleConnector( new InetSocketAddress(127.0.0.1, 11111), example, , ); int batchSize 1000; // 每次获取的消息数量 try { connector.connect(); // 连接 connector.subscribe(.*\\..*); // 订阅过滤条件可以和Server配置不同进行二次过滤 connector.rollback(); // 回滚到未ack的位置如果是第一次则从当前binlog开始 while (true) { Message message connector.getWithoutAck(batchSize); // 获取指定数量的数据不自动确认 long batchId message.getId(); if (batchId ! -1 !message.getEntries().isEmpty()) { // 2. 解析Message中的Entry printEntries(message.getEntries()); connector.ack(batchId); // 提交确认表示成功消费。如果处理失败可以调用rollback(batchId)回滚。 } else { Thread.sleep(1000); // 没有数据休息一下 } } } catch (Exception e) { e.printStackTrace(); } finally { connector.disconnect(); } } private static void printEntries(ListCanalEntry.Entry entries) { for (CanalEntry.Entry entry : entries) { if (entry.getEntryType() EntryType.ROWDATA) { RowChange rowChange; try { rowChange RowChange.parseFrom(entry.getStoreValue()); } catch (Exception e) { throw new RuntimeException(解析RowChange错误, e); } EventType eventType rowChange.getEventType(); String tableName entry.getHeader().getTableName(); System.out.println( 监听到表[ tableName ]的 eventType 事件 ); for (RowData rowData : rowChange.getRowDatasList()) { if (eventType EventType.DELETE) { // 删除操作rowData.getBeforeColumnsList()包含被删除行的所有旧值 printColumns(rowData.getBeforeColumnsList()); } else if (eventType EventType.INSERT) { // 插入操作rowData.getAfterColumnsList()包含新插入行的所有值 printColumns(rowData.getAfterColumnsList()); } else if (eventType EventType.UPDATE) { // 更新操作Before是旧值After是新值 System.out.println(--- 更新前 ---); printColumns(rowData.getBeforeColumnsList()); System.out.println(--- 更新后 ---); printColumns(rowData.getAfterColumnsList()); } } } } } private static void printColumns(ListColumn columns) { for (Column column : columns) { System.out.println(column.getName() : column.getValue() (更新 column.getUpdated() )); } } }这段代码做了几件事连接Canal、订阅数据、循环拉取消息、解析出具体的增删改事件和字段数据。关键点在于ack()和rollback()这实现了客户端的消费确认机制保证数据至少被消费一次at-least-once。注意事项客户端的消费速度必须跟上数据库的变更速度否则消息会在Canal Server的内存中堆积最终导致内存溢出OOM。你需要监控客户端的消费延迟可以从Message的getId大致判断它是一个递增的批次ID也可以从Canal Server的日志或管理端口查看。如果延迟变大需要考虑优化客户端处理逻辑或者增加客户端消费能力如多线程处理。4.2 使用Canal Adapter进行便捷数据同步如果你不想写代码只是想简单快速地把MySQL数据同步到其他存储比如Elasticsearch、Redis或者另一个MySQL那么Canal Adapter是更好的选择。它是一个开箱即用的客户端通过配置文件就能完成同步任务。部署Adapter同样从Release页面下载canal.adapter-1.1.7.tar.gz并解压。修改其配置文件conf/application.yml。canal.conf: mode: tcp # 连接模式与Server的canal.serverMode对应 canalServerHost: 127.0.0.1:11111 # Canal Server地址 batchSize: 500 syncBatchSize: 1000 retries: 0 timeout: srcDataSources: # 源数据源就是Canal监听的MySQL配置可以多个 defaultDS: url: jdbc:mysql://127.0.0.1:3306/your_database?useUnicodetrue username: canal password: canal canalAdapters: - instance: example # Canal实例名 groups: - groupId: g1 outerAdapters: - name: logger # 第一个适配器日志输出用于测试 - name: es7 # 第二个适配器同步到Elasticsearch 7.x hosts: 127.0.0.1:9200 properties: mode: rest security.auth: user:password # 指定具体的同步映射配置文件 - name: rdb key: mysql_es_user # 自定义key用于在mapping配置中引用此数据源 properties: jdbc.driverClassName: com.mysql.jdbc.Driver jdbc.url: jdbc:mysql://127.0.0.1:3306/your_database jdbc.username: canal jdbc.password: canal在conf/rdb目录下创建具体的表映射配置文件例如user.ymldataSourceKey: defaultDS # 对应application.yml中的srcDataSources destination: example # Canal实例名 groupId: g1 outerAdapterKey: mysql_es_user # 对应application.yml中rdb适配器的key concurrent: true dbMapping: database: your_database # 源数据库 table: user # 源表 targetTable: user # 目标表对于ES是索引名 targetPk: id # 主键 mapAll: true # 映射所有字段 # insertOp: I # 操作映射默认即可 # updateOp: U # deleteOp: D启动Adapterbin/startup.sh。查看日志logs/adapter/adapter.log如果看到成功订阅并开始同步的日志就说明配置成功了。Adapter将自动监听Canal Server的数据并根据你的映射配置将数据写入到指定的目标数据源。这种方式极大地简化了ETL流程。5. 生产环境部署进阶与高可用考量单机部署只能用于测试和学习。生产环境必须考虑高可用和可扩展性。5.1 集群化部署基于ZooKeeperCanal支持将多个Server节点组成集群通过ZooKeeper进行实例管理和故障转移。这样当一个Canal Server节点宕机时ZooKeeper会将其负责的实例切换到其他存活节点上保证服务不中断。配置步骤概要部署ZooKeeper集群。修改每个Canal Server节点的canal.propertiescanal.zkServers zk1:2181,zk2:2181,zk3:2181 canal.instance.global.mode spring canal.instance.global.lazy false canal.instance.global.manager.address ${canal.zkServers} canal.instance.global.spring.xml classpath:spring/default-instance.xml将持久化模式从file改为spring并将管理器地址指向ZooKeeper。修改实例配置文件instance.properties关键一步指定该实例由哪个canalId即Server节点来运行。如果不指定则由ZooKeeper动态分配。canal.instance.mysql.slaveId 1234 # 需全局唯一 # 可选固定该实例到某个canal主机 # canal.instance.canal.manager.canalId canal-server-hostname:11111启动所有Canal Server节点它们会自动向ZooKeeper注册。你可以通过Calan AdminCanal的管理界面来查看集群状态和手动进行主备切换。5.2 性能调优与监控网络与缓冲区在canal.properties中可以调整canal.instance.network.receiveBufferSize和canal.instance.network.sendBufferSize来优化网络吞吐。调整canal.instance.memory.buffer.size和canal.instance.memory.batch.mode可以平衡内存使用和消费延迟。解析线程数对于有大量表或频繁变更的场景可以增加解析线程数canal.instance.parser.parallelThreadSize默认为1。监控Canal暴露了JMX指标。你可以通过canal.metrics.pull.port配置端口然后使用Prometheus的JMX Exporter或直接通过JMX客户端来采集监控数据如消费延迟、内存使用率、事件处理速率等并配置告警。日志切割Canal使用logback可以在conf/logback.xml中配置按天或按大小切割日志避免日志文件过大。5.3 与消息队列Kafka/RocketMQ集成在超大规模数据同步场景下直接使用TCP模式可能会因为客户端消费能力不均或网络问题导致Server不稳定。更优雅的方案是让Canal Server将解析后的数据直接投递到Kafka或RocketMQ让消息队列来承担削峰填谷和解耦的责任。配置非常简单只需修改canal.propertiescanal.serverMode kafka # 或 rocketmq # Kafka配置示例 kafka.bootstrap.servers 127.0.0.1:9092 kafka.acks all kafka.compression.type snappy kafka.batch.size 16384 kafka.linger.ms 100同时在instance.properties中配置投递的Topic和分区规则。这样Canal Server就变成了一个可靠的生产者客户端则从Kafka消费实现了存储与消费的完全解耦系统的伸缩性和容错性大大增强。6. 常见问题排查与实战踩坑记录即使按照步骤来也难免会遇到问题。这里我总结几个最常见的坑和排查思路。问题1Canal启动失败日志显示“Address already in use”原因11111端口被占用。解决netstat -tlnp | grep 11111找出占用进程并停止或者修改canal.properties中的canal.port为其他端口。问题2Canal连接MySQL失败日志显示“Access denied for user ‘canal‘”原因权限不足或密码错误。解决确认MySQL中‘canal‘‘%‘用户已创建且密码正确。确认Canal服务器IP是否在MySQL的授权范围内%代表所有主机生产环境建议指定IP。在Canal服务器上用mysql命令行工具测试连接mysql -h your_mysql_host -ucanal -p。问题3客户端连接Canal Server失败原因网络不通、防火墙、或Canal Server未正常启动。解决telnet canal_server_ip 11111测试端口连通性。检查Canal Server日志canal.log和example.log是否有错误。确认客户端代码中连接地址和实例名destination是否正确。问题4能连接但消费不到数据原因MySQL的binlog格式不是ROW。订阅过滤规则filter.regex配置错误没有覆盖到你操作的表。客户端订阅的过滤条件与Server端或实际操作不匹配。位点消费进度问题可能已经消费过了最新的数据。解决确认binlog_format ROW。在MySQL中执行一个简单的INSERT/UPDATE操作观察Canal的example.log是否有DML相关的日志。如果有说明Server端解析到了问题在客户端如果没有检查Server配置。在客户端代码中尝试先执行connector.subscribe(“.*\\..*”)订阅所有看能否收到数据。尝试在客户端调用connector.rollback()回滚到未ack的位置或者使用connector.get(batchSize)自动ack看能否立即获取到刚产生的数据。问题5同步延迟越来越大原因客户端消费速度跟不上数据库变更速度。解决优化客户端处理逻辑比如采用批量处理、异步处理。增加客户端消费并行度例如启动多个客户端消费同一个实例的不同分区需结合Kafka模式或调整客户端参数。检查网络带宽和Canal Server所在机器的负载CPU、IO。如果使用的是Adapter同步到ES等目标检查目标端的写入性能是否成为瓶颈。问题6同步的数据出现乱码原因字符集不一致。解决确保MySQL数据库、表、字段的字符集Canal Server配置中的connectionCharset以及客户端或Adapter处理数据时的字符集如JVM默认编码、ES索引mapping定义全部统一推荐使用UTF-8。部署和运维Calan的过程就是一个不断与这些细节打交道的过程。我的经验是一定要把日志级别调到DEBUG修改logback.xml来排查复杂问题同时做好全方位的监控。一旦稳定运行起来Canal会成为你数据架构中非常可靠的一环。
返回列表