尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Canal实战指南:基于MySQL Binlog的实时数据同步原理与部署

Canal实战指南:基于MySQL Binlog的实时数据同步原理与部署 1. 项目概述为什么我们需要Canal如果你负责过数据相关的项目大概率遇到过这样的场景业务数据库里一张表的数据变更了下游的缓存、搜索引擎或者数据分析系统需要立刻感知到以便同步更新自己的数据。最直接的做法是让业务代码在写数据库的同时再写一遍这些下游系统但这会带来代码耦合、事务难以保证、性能下降等一系列问题。另一种常见的思路是定时去数据库里“扫表”但这又存在延迟高、对数据库有查询压力、可能漏掉变更的弊端。Canal这个由阿里巴巴开源的项目就是为了优雅地解决这类“数据库变更实时捕获与同步”的需求而生的。它的名字取自英文“Channel”寓意着数据流动的管道。简单来说Canal伪装成MySQL数据库的一个从库Slave向主库Master发起数据同步请求。主库会将自己的二进制日志binlog推送给CanalCanal解析这些日志就能得到数据库中每一行数据的增、删、改的详细信息然后将这些变更事件以结构化的消息例如JSON推送出去。下游的任何系统无论是Redis、Elasticsearch还是Kafka、Flink只需要订阅这些消息就能实现近乎实时的数据同步。这套机制的核心价值在于解耦和实时性。业务代码只需关心把数据写入MySQL无需感知下游有多少个消费者。而下游系统通过监听一个统一的事件流来获取变更延迟可以控制在毫秒级。这对于构建现代数据架构如缓存更新、搜索索引构建、实时数仓、跨系统数据一致性保障等场景是一个基石性的组件。我经历过从业务代码双写到引入Canal的改造最大的感受是系统边界清晰了数据流的维护成本直线下降。2. Canal核心原理与架构拆解要玩转Canal不能只停留在“会用”的层面必须理解其内部的工作原理和组件构成。这能帮助你在部署、配置和排查问题时快速定位根因。2.1 基于MySQL主从复制的实现原理Canal的技术灵感直接来源于MySQL自身的主从复制Replication机制。理解这一点是理解Canal所有行为的基础。MySQL主从复制的过程大致如下主库Master将数据变更DML、DDL记录到二进制日志Binary Log简称binlog中。从库Slave的I/O线程会连接到主库请求读取binlog。主库的Binlog Dump线程会根据从库的请求将binlog内容推送给从库。从库的I/O线程将接收到的binlog数据写入本地的中继日志Relay Log。从库的SQL线程读取中继日志并重放Replay其中的SQL事件从而使从库的数据与主库保持一致。Canal的核心创新在于它把自己伪装成了一个MySQL从库。它实现了MySQL从库与主库交互的协议可以向主库发送COM_BINLOG_DUMP命令。这样主库就会像对待一个真正的从库一样持续地将binlog推送给Canal。注意这意味着你的MySQL必须开启binlog并且格式设置为ROW模式。STATEMENT或MIXED模式无法提供精确到行级别的变更前、后的数据Canal的解析能力会大打折扣甚至无法工作。Canal接收到binlog流之后并不会像真正的从库那样去执行SQL。它的工作是对这些二进制的binlog进行解析还原出每一个变更事件Event的详细信息包括变更类型INSERT、UPDATE、DELETE。数据库名、表名。变更前数据Before Image对于UPDATE和DELETE至关重要。变更后数据After Image对于INSERT和UPDATE至关重要。变更时间戳、事务ID等元信息。解析完成后Canal将这些信息封装成更易处理的结构通常是Protobuf格式并存储或转发。整个过程中Canal对主库的影响与一个只读的从库完全一致只有非常轻微的网络和IO消耗。2.2 Canal服务端核心组件解析一个标准的Canal服务端部署包含以下几个关键部分理解它们有助于进行监控和调优Server代表一个Canal运行实例。一个Server可以包含多个Instance。Instance这是实际的工作单元。一个Instance对应一个数据源一个MySQL数据库实例负责该数据源的binlog订阅、解析和投递。通常我们为一个业务数据库创建一个Instance。EventParser核心的解析器。它负责连接MySQL模拟从库协议拉取binlog并调用EventTransactionBuffer进行事务合并和排序最终将解析结果传递给EventSink。EventSink事件的“过滤器”和“连接器”。它负责对Parser解析出的事件进行过滤例如根据配置的表名规则和归并将同一个数据库的多个表变更归并到一个数据队列然后投递给EventStore。EventStore的设计模式使得Parser和Store之间可以解耦。EventStore事件的存储层。目前主要实现是MemoryEventStore即内存存储。它采用类似RingBuffer的机制分为Put和Get两个操作序列。Parser是Producer向Store中写入事件Canal客户端是Consumer从Store中获取事件。这个内存缓冲区的大小是可配置的它直接决定了Canal能缓冲多少未消费的变更事件是影响稳定性的关键参数。MetaManager元数据管理器。它负责记录Canal客户端消费的进度类比Kafka的Consumer Offset。例如客户端A消费到了binlog filemysql-bin.000001的position107这个位置这个信息就需要被持久化以便客户端重启后能从断点继续消费避免数据重复或丢失。支持多种存储方式如内存、文件、ZooKeeper等。------------------- ------------------- ------------------- | MySQL Master |-----| Canal Instance |-----| Client (Java) | | (Binlog) |-----| (Parser/Sink/ |-----| (订阅消费) | ------------------- | Store) | ------------------- ------------------- | 元数据 v ------------------- | MetaManager | | (ZK/File/Memory) | -------------------2.3 部署模式Standalone vs. HA根据可用性要求Canal有两种主要的部署模式Standalone单机模式最简单直接的部署方式。一个Canal服务进程包含一个或多个Instance。配置简单启动快速适合开发、测试环境或对可用性要求不高的生产场景。其缺点是存在单点故障如果Canal服务宕机整个数据同步链路就会中断。HA高可用模式为了解决单点问题Canal支持基于ZooKeeper或其它协调服务的HA部署。其核心思想是为同一个Instance启动多个Canal服务器节点例如canal-server-1,canal-server-2它们共同竞争一个在ZooKeeper上创建的临时节点EPHEMERAL锁。获得锁的节点成为Active节点承担实际的binlog解析和投递工作其他节点作为Standby处于休眠状态持续监控锁的状态。故障转移当Active节点宕机或与ZK失联其持有的临时锁会自动释放。Standby节点中会有一个迅速抢到锁升级为新的Active节点并从MetaManager中读取上一个Active节点记录的消费位点接着从这个位置开始拉取binlog从而实现服务的无缝切换。客户端适配HA模式下的Canal客户端如CanalConnector在连接时会首先从ZooKeeper获取当前Active服务器的地址然后与之建立连接。当发生故障转移后客户端能感知到Active节点的变化并自动重连到新的服务器。实操心得在生产环境除非数据同步链路可接受短暂中断否则强烈建议使用HA模式。部署时确保多个Canal服务器节点的canal.properties中canal.zkServers配置指向相同的ZooKeeper集群并且每个Instance的canal.instance.global.spring.xml配置主要是数据库连接信息完全一致。我曾因为两个节点配置文件里数据库密码有一个字符不同导致切换后Standby节点无法连接MySQLHA形同虚设。3. 从零开始搭建与配置Canal服务端理论清楚了我们动手搭建一个Canal服务端。这里以Standalone模式为例HA模式在此基础上增加ZK配置即可。3.1 环境准备与依赖检查首先确保你的基础环境就绪JavaCanal 1.1.x版本需要JDK 1.8或以上。运行java -version确认。MySQL版本5.6, 5.7, 8.0均可。这是数据源头必须满足以下条件开启binlog编辑MySQL配置文件如/etc/my.cnf或/etc/mysql/mysql.conf.d/mysqld.cnf确保有以下配置[mysqld] # 启用binlog并设置文件名前缀 log-binmysql-bin # 设置binlog格式为ROW这是Canal工作的必要条件 binlog-formatROW # 为MySQL集群中的每台服务器设置唯一的ID对于单机任意正整数即可 server-id1 # 8.0版本可能需要显式指定binlog行镜像为FULL确保记录变更前后所有列的值 binlog_row_imageFULL修改后重启MySQL服务。创建Canal专用账号Canal需要以从库身份连接MySQL因此需要一个具有REPLICATION SLAVE和REPLICATION CLIENT权限的账号。CREATE USER canal% IDENTIFIED BY canal_password; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; FLUSH PRIVILEGES;注意生产环境请将%替换为Canal服务器具体的IP地址并设置强密码。3.2 Canal服务端下载与基础配置下载发布包从Canal的GitHub Release页面下载最新稳定版如canal.deployer-1.1.7.tar.gz。解压与目录结构tar -zxvf canal.deployer-1.1.7.tar.gz -C /opt/ cd /opt/canal关键目录说明conf/配置文件目录。canal.propertiesCanal Server级别的全局配置。example/一个名为example的Instance配置模板目录。logs/日志目录。lib/依赖库。修改全局配置(conf/canal.properties) 这个文件配置Canal Server的运行参数。对于初次使用大部分默认值即可。重点关注以下几项# Canal Server的工作模式tcp表示提供原生Socket服务kafka/rocketmq表示解析后直接投递到消息队列。我们先从tcp开始。 canal.serverMode tcp # 每个Instance的端口偏移量。example实例的端口 canal.port(默认11111) 1 11112 canal.port 11111 # Instance列表多个用逗号分隔。这里定义了一个叫example的实例。 canal.destinations example # 可配置conf根目录下的哪些目录作为Instance配置目录默认就是example。 canal.conf.dir ../conf # 自动扫描Instance配置变化的间隔毫秒 canal.auto.scan true # 内存存储的批次大小影响客户端一次获取的消息条数 canal.instance.memory.batch.size 1000修改Instance配置(conf/example/instance.properties) 这个文件配置具体连接哪个MySQL以及解析行为。# 数据源地址 canal.instance.master.address127.0.0.1:3306 # 数据库账号密码刚才创建的 canal.instance.dbUsernamecanal canal.instance.dbPasswordcanal_password # 字符集 canal.instance.connectionCharsetUTF-8 # 需要订阅的库表过滤规则默认是.*\\..*即所有库所有表。 # 例如只订阅test库的所有表test\\..* # 只订阅test库的user表test.user canal.instance.filter.regex.*\\..* # 排除某些表的规则默认为空。 # canal.instance.filter.black.regex重要提示canal.instance.filter.regex是控制数据过滤的第一道关口正确设置可以极大减少不必要的网络传输和解析开销。建议根据业务需求精确配置。3.3 启动服务与日志查看启动Canal Servercd /opt/canal sh bin/startup.shWindows环境下使用startup.bat。检查启动日志tail -f logs/canal/canal.log看到## the canal server is running now ......即表示Server启动成功。查看Instance日志tail -f logs/example/example.log重点关注是否有连接MySQL成功、binlog位置加载成功等信息。如果看到prepare to find start position ...和find start position successfully ...说明Instance启动并准备就绪。如果启动失败常见原因有MySQL未开启binlog或格式不是ROW。数据库账号权限不足。网络不通或MySQL的bind-address配置为127.0.0.1导致外部无法连接。端口冲突。4. 开发Canal客户端订阅与处理数据变更服务端跑起来了现在我们需要一个客户端来消费变更数据。Canal提供了多种客户端接入方式最常用的是使用其官方Java客户端。4.1 引入Java客户端依赖创建一个Maven项目添加依赖dependency groupIdcom.alibaba.otter/groupId artifactIdcanal.client/artifactId version1.1.7/version !-- 与服务端版本保持一致 -- /dependency如果你使用Spring Boot也可以直接引入。4.2 编写一个简单的客户端示例下面是一个最基础的客户端代码它连接到Canal服务端订阅example实例并持续打印收到的数据变更。import com.alibaba.otter.canal.client.CanalConnector; import com.alibaba.otter.canal.client.CanalConnectors; import com.alibaba.otter.canal.protocol.CanalEntry.*; import com.alibaba.otter.canal.protocol.Message; import java.net.InetSocketAddress; import java.util.List; public class SimpleCanalClient { public static void main(String[] args) { // 1. 创建连接器。这里连接Standalone模式的服务端。 // 如果是HA模式应使用 CanalConnectors.newClusterConnector(zkServers, destination, username, password) String destination example; // 对应 instance.properties 配置的目录名 CanalConnector connector CanalConnectors.newSingleConnector( new InetSocketAddress(127.0.0.1, 11111), // Canal Server地址 destination, , // 用户名默认空 // 密码默认空 ); final int batchSize 1000; // 每次获取的消息数量 try { // 2. 建立连接 connector.connect(); // 3. 订阅过滤条件。这里订阅所有表也可以指定库表 connector.subscribe(test.user); connector.subscribe(.*\\..*); // 4. 回滚到未确认的位置如果是重启会从上次ack的位置开始 connector.rollback(); while (true) { // 5. 获取指定数量的消息 Message message connector.getWithoutAck(batchSize); long batchId message.getId(); int size message.getEntries().size(); if (batchId -1 || size 0) { // 没有数据稍作休息 Thread.sleep(1000); } else { // 6. 处理消息 printEntries(message.getEntries()); } // 7. 确认消息。确认后Canal Server会删除这部分消息并更新消费位点。 connector.ack(batchId); // 如果处理失败可以回滚 connector.rollback(batchId); } } catch (Exception e) { e.printStackTrace(); } finally { // 8. 断开连接 connector.disconnect(); } } private static void printEntries(ListCanalEntry.Entry entries) { for (CanalEntry.Entry entry : entries) { // 只处理事务开始、数据变更、事务结束这三种类型的Entry if (entry.getEntryType() EntryType.TRANSACTIONBEGIN || entry.getEntryType() EntryType.TRANSACTIONEND) { continue; } RowChange rowChange; try { rowChange RowChange.parseFrom(entry.getStoreValue()); } catch (Exception e) { throw new RuntimeException(解析RowChange出错, e); } EventType eventType rowChange.getEventType(); // 打印头部信息 System.out.println(String.format( binlog[%s:%s], db[%s], table[%s], type[%s], entry.getHeader().getLogfileName(), entry.getHeader().getLogfileOffset(), entry.getHeader().getSchemaName(), entry.getHeader().getTableName(), eventType)); // 处理每一行数据变更 for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { if (eventType EventType.DELETE) { // DELETE操作打印变更前的数据 printColumns(rowData.getBeforeColumnsList()); } else if (eventType EventType.INSERT) { // INSERT操作打印变更后的数据 printColumns(rowData.getAfterColumnsList()); } else if (eventType EventType.UPDATE) { // UPDATE操作打印变更前和变更后的数据 System.out.println(------- 更新前); printColumns(rowData.getBeforeColumnsList()); System.out.println(------- 更新后); printColumns(rowData.getAfterColumnsList()); } } } } private static void printColumns(ListCanalEntry.Column columns) { for (CanalEntry.Column column : columns) { System.out.println(column.getName() : column.getValue() update column.getUpdated()); } } }运行这个客户端然后在对应的MySQL数据库里执行一些INSERT、UPDATE、DELETE操作你就能在控制台看到实时的变更输出了。4.3 客户端核心逻辑与最佳实践上面的示例展示了基本流程但在生产环境中我们需要考虑更多连接管理客户端应具备重连机制。网络波动或Canal服务重启可能导致连接断开客户端需要能够自动检测并重新连接、重新订阅。消息确认ACK机制connector.getWithoutAck()获取消息后必须调用connector.ack(batchId)进行确认。只有确认后Canal Server才会将这批消息标记为已消费并更新消费位点写入MetaManager。如果处理失败应调用connector.rollback(batchId)这样下次获取还会从这批消息开始。务必确保业务处理成功后再ACK否则可能导致数据丢失。批量处理getWithoutAck的batchSize参数需要根据业务处理速度和内存情况权衡。太小会增加网络和ACK开销太大可能导致单次处理时间过长内存占用高且一旦失败回滚的数据量也大。通常设置在几百到几千之间。消费位点管理Canal客户端默认会将消费位点定期持久化。对于Java客户端如果使用CanalConnectors.newClusterConnector连接HA模式位点信息会保存在ZooKeeper上。确保你的客户端有权限读写ZK上的相应路径。数据解析与转换示例中直接打印了原始数据。实际应用中你需要将RowData转换成你的业务对象DTO或者直接组装成发送给下游系统如ES、Redis的命令。注意处理各种数据类型如日期时间、二进制、JSON等的转换。异常处理与监控务必对消息解析、业务处理等环节做好异常捕获和日志记录。同时可以监控客户端消费的延迟通过对比当前时间与binlog中的时间戳延迟过高可能意味着下游处理能力不足。实操心得在消费逻辑中我强烈建议将消息解析和业务处理两个步骤解耦。解析层只负责将Canal的Protobuf消息转换成简单的POJO事件包含库、表、操作类型、变更前后数据等。业务处理层监听这些POJO事件。这样做的好处是当Canal协议升级或你想更换消息来源比如换成Debezium时只需改动解析层业务代码几乎不受影响。5. 高级配置与性能调优指南当数据量增大或业务要求变高时默认配置可能无法满足需求。以下是一些关键的性能调优和高级配置点。5.1 服务端关键参数调优配置文件conf/canal.properties:参数默认值说明与调优建议canal.instance.memory.batch.size1000内存存储批次大小。客户端每次get请求能获取的最大消息条数。增大此值可以提高吞吐减少网络交互但会增加客户端单次处理的内存压力和延迟。根据客户端消费能力调整。canal.instance.memory.buffer.size16384内存存储缓冲区大小。这是MemoryEventStore环形缓冲区的大小单位是事件个数。这是最重要的参数之一。如果生产者Parser速度持续快于消费者缓冲区会被写满导致Parser线程阻塞进而可能使MySQL主库的Binlog Dump线程被阻塞。必须根据业务峰值流量设置足够大的值。计算公式可粗略估算为(峰值TPS * 事务平均事件数 * 预计最大消费延迟秒数)。例如峰值1000 TPS平均一个事务10个事件允许延迟5秒则至少需要1000*10*550000。建议设置为2的幂。canal.instance.transaction.size1024事务合并批次大小。Parser在解析binlog时会尝试将多个小事务合并后再投递给Sink。增大此值可以减少Store的写入次数提升吞吐但会略微增加延迟。在事务非常小的场景下如单行插入可以适当调大。canal.instance.parser.paralleltrue是否启用并行解析。对于多核机器开启并行解析可以提升解析binlog的效率。建议开启。canal.instance.parser.parallelThreadSize自动计算并行解析的线程数。默认根据CPU核数计算通常无需调整。配置文件conf/example/instance.properties:参数默认值说明与调优建议canal.instance.connectionCharsetUTF-8连接MySQL的字符集必须与数据库一致否则可能出现乱码。canal.instance.filter.regex.*\\..*白名单过滤。务必按需设置只订阅必要的库表。过滤在Parser解析后、Sink处理前进行减少不必要的数据流动。支持Perl风格正则。canal.instance.filter.black.regex空黑名单过滤。优先级低于白名单。canal.instance.fallbackIntervalInSeconds60降级间隔。当发生网络异常等问题时Canal会尝试降级到定时从MySQL拉取数据的方式。此配置设置定时拉取的间隔。一般保持默认。canal.instance.detecting.enablefalse是否开启心跳检测。开启后Canal会定期向MySQL执行SELECT 1以保持连接活性。在防火墙会话时间较短的环境中可以开启。canal.instance.detecting.sqlselect 1心跳检测SQL。5.2 选择合适的消费模式除了上面演示的TCP直连模式canal.serverMode tcpCanal还支持将解析后的数据直接投递到消息队列这更适合大规模、多消费者的生产环境。Kafka/RocketMQ模式(canal.serverMode kafka或rocketmq)优势解耦更彻底Canal Server只负责解析和投递到MQ消费压力由MQ集群承担。支持多消费者下游多个系统可以独立消费同一份数据。消息堆积与回溯利用MQ的消息堆积能力可以应对下游消费能力不足的情况也支持按时间偏移量重新消费。高可用依托于Kafka/RocketMQ自身的高可用机制。配置需要在canal.properties中配置MQ的地址、Topic等参数并在instance.properties中配置每个Instance对应的具体Topic。注意事项消息格式需要约定好。Canal投递到MQ的消息体默认仍是Protobuf序列化后的字节下游消费者需要相应的反序列化代码。也可以配置为JSON等格式。模式选择建议开发测试、简单同步任务TCP直连模式简单快捷。生产环境、单一消费者如果下游只有一个系统且对可靠性要求不是极端高TCP直连HA模式也是可选的。生产环境、多消费者、大数据量无脑选择Kafka/RocketMQ模式。这是目前最主流、最稳健的架构。5.3 监控与告警没有监控的系统就是在“裸奔”。Canal的监控主要关注以下几点服务状态Canal进程是否存活。可以通过简单的进程检查或健康检查接口Canal自带一个简单的管理端口默认11110来监控。消费延迟这是最重要的业务指标。延迟计算公式当前时间 - 最后一条已消费binlog事件的时间戳。可以通过解析Canal的消费位点binlog file position连接到MySQL执行SHOW BINARY LOGS和SHOW BINLOG EVENTS来估算或者从Canal的JMX指标中获取如果开启。延迟持续增大说明下游消费能力不足或Canal/MQ有瓶颈。内存缓冲区使用率监控MemoryEventStore的put和get序列号之差。如果这个差值持续接近buffer.size说明缓冲区即将写满Parser可能被阻塞需要紧急扩容缓冲区或提升下游消费速度。解析与投递速率监控Canal每秒解析和投递的事件数TPS。可以与MySQL的写操作TPS对比正常情况下应该接近。错误日志监控logs/目录下的error日志及时发现连接异常、解析错误等问题。可以将这些指标通过JMX或自定义日志输出接入到Prometheus、Zabbix等监控系统并设置合理的告警阈值。6. 生产环境常见问题与排查实录即使配置得当在生产环境中运行Canal仍可能遇到各种问题。下面是我在实践中总结的一些典型问题及其排查思路。6.1 数据丢失或重复消费这是最令人头疼的问题。可能原因及排查客户端ACK机制使用不当这是最常见的原因。业务代码处理消息时发生异常但依然执行了ack()导致消息被确认但实际上业务并未成功处理。务必确保在业务逻辑完全成功后再调用ack()。建议使用try-catch-finally结构在finally中根据处理结果决定ack还是rollback。消费位点未正确持久化检查MetaManager的配置。如果是文件模式确保canal.instance.data.dir目录有写入权限且磁盘空间充足。如果是ZK模式检查ZK连接是否稳定客户端是否有写权限。位点丢失会导致客户端重启后从错误的通常是更早的位置开始消费造成数据重复。Canal Server缓冲区溢出如果canal.instance.memory.buffer.size设置过小而消费速度长期慢于生产速度缓冲区写满后Parser线程会被阻塞。在TCP模式下这可能导致MySQL主库的Binlog Dump线程也被阻塞极端情况下可能触发MySQL的写超时甚至导致Canal丢失尚未解析的binlog事件如果binlog文件被Purge。必须根据流量评估并设置足够大的缓冲区并监控其使用率。网络分区或长时间GC在HA模式下如果Active节点发生长时间GC或网络分区可能导致ZK会话超时锁被释放Standby节点切换为Active。但原Active节点GC恢复或网络恢复后可能并未感知到自己已不是Active会继续解析binlog并投递造成双写导致下游数据重复。虽然Canal有机制避免但在极端情况下仍可能发生。需要监控服务器GC情况和网络健康。6.2 同步延迟高下游系统发现数据更新很久才同步过来。可能原因及排查下游消费能力不足这是最主要的原因。检查下游消费者你的业务程序的处理逻辑是否太慢是否有耗时的IO操作、同步RPC调用等。考虑优化消费逻辑或增加消费者数量如果是MQ模式。Canal Server或MQ瓶颈CPU/IO监控Canal Server所在机器的CPU和磁盘IO使用率。解析binlog是CPU密集型操作。网络带宽如果单行数据量很大如包含TEXT、BLOB字段同步大量数据可能占满网络带宽。MQ堆积如果是MQ模式检查Kafka/RocketMQ是否有消息堆积。可能是MQ集群性能问题或Topic分区数太少导致消费并发度不够。不合理的批量大小客户端batchSize设置过大单次处理耗时过长虽然吞吐可能高但延迟也会增加。可以适当调小batchSize牺牲一些吞吐来换取更低的延迟。MySQL主库压力大如果MySQL主库本身负载很高生成binlog和响应Canal作为从库的拉取请求也会变慢。6.3 连接异常与解析错误Connect to mysql server ... failed检查MySQL服务是否正常网络是否连通。检查canal.instance.dbUsername和canal.instance.dbPassword是否正确。检查MySQL用户canal的权限和主机限制canal%。检查MySQL的max_allowed_packet参数是否过小导致大的binlog事件无法传输。Could not find first log file name in the binary log index fileCanal启动时会从MetaManager中读取上次消费的位点。如果位点信息丢失或指定的binlog文件已被Purge清除就会报此错误。解决可以手动修改Instance的Meta信息。对于文件模式编辑conf/example/meta.dat文件谨慎操作。更常见的做法是如果允许从头同步可以删除这个meta文件Canal会从当前最新的binlog位置开始消费。但这会导致数据重复另一种是指定一个较新的、存在的binlog文件位置。在instance.properties中配置# 指定起始位置覆盖meta.dat中的记录 canal.instance.master.journal.namemysql-bin.000123 canal.instance.master.position456789 canal.instance.master.timestamp1609459200000 # 三者任选其一即可优先级journal.name position timestamp解析DML语句出错通常是因为表结构发生了变化ALTER TABLE而Canal缓存的表元信息Schema是旧的。Canal会在第一次解析表时缓存其结构后续如果表结构变更需要重启Canal Instance以重新获取最新的Schema。对于在线业务频繁重启不可接受。可以考虑启用Canal的表结构自动刷新功能需要MySQL开启binlog_rows_query_log_events或者使用一些开源的管理工具来动态刷新。6.4 与MySQL 8.0的兼容性问题MySQL 8.0默认使用了新的密码认证插件caching_sha2_password而旧版本的Canal客户端驱动可能不支持。解决方案推荐升级Canal使用较新版本的Canal如1.1.6其内置的MySQL驱动已支持新的认证方式。修改MySQL用户认证插件如果可行ALTER USER canal% IDENTIFIED WITH mysql_native_password BY canal_password; FLUSH PRIVILEGES;这将用户的认证方式改回旧的mysql_native_password。踩过这些坑之后我的体会是对于Canal这类数据同步组件监控和预警必须走在前面。不能等到业务方投诉数据没同步才发现问题。要像对待业务核心服务一样为它建立完善的监控仪表盘关注延迟、缓冲区、错误率等核心指标并设置有效的告警规则。同时任何对Canal配置的变更尤其是缓冲区大小、过滤规则等都应在测试环境充分验证后再上线。
返回列表