ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Canal 全量+增量数据同步:Adapter/Client 集成方案与全量初始化坑点

Canal 全量+增量数据同步:Adapter/Client 集成方案与全量初始化坑点 Canal 全量增量数据同步Adapter/Client 集成方案与全量初始化坑点1. Canal 基本原理与架构Canal 是阿里巴巴开源的一款基于 MySQL 数据库增量日志解析的组件它通过监听 MySQL 的 binlog 日志实现对数据库变更的捕获和分发。Canal 的核心架构包括Canal Server: 负责连接 MySQL 数据库解析 binlog并将解析后的数据推送给客户端Canal Client: 接收 Server 推送的数据进行消费和处理Adapter: 提供现成的适配器支持多种目标存储全量增量同步的基本思想是首先通过某种机制将源数据库的全量数据同步到目标数据库然后通过 Canal 持续同步增量变更实现数据的最终一致性。2. Adapter 集成方案详解Canal 提供了多种适配器如 ElasticSearch、HBase、Redis 等可以直接配置使用无需编写自定义代码。2.1 Adapter 集成步骤下载并部署 Canal Serverbashwget https://github.com/alibaba/canal/releases/download/canal-1.1.5/canal.deployer-1.1.5.tar.gztar -zxvf canal.deployer-1.1.5.tar.gz配置 Canal Server修改conf/example/instance.properties文件配置 MySQL 连接信息propertiescanal.dbUsername canalcanal.dbPassword canalcanal.dbHostname 127.0.0.1canal.dbPort 3306下载并部署 Canal Adapterbashwget https://github.com/alibaba/canal/releases/download/canal-adapter-1.1.5/canal.adapter-1.1.5.tar.gztar -zxvf canal.adapter-1.1.5.tar.gz配置 Adapter修改conf/application.yml文件配置目标存储信息yamlserver:port: 8081spring:jackson:date-format: yyyy-MM-dd HH:mm:sstime-zone: GMT8canal.conf:mode: rabbitmqmqServers: 127.0.0.1:5672flatMessage: truezooKeeperHosts:batchMode: falsesyncBatchSize: 1024retries: 0timeout:srcDataSources:defaultDS:url: jdbc:mysql://127.0.0.1:3306/mydb?useUnicodetruecharacterEncodingUTF-8username: rootpassword: rootcanalAdapters:instance: examplegroups:groupId: g1outerAdapter:key: mysql1type: mysqlhosts: 127.0.0.1:3306username: rootpassword: rootdatabase: mydb_target启动服务bash# 启动 Canal Serversh bin/startup.sh# 启动 Canal Adaptersh bin/adapter/startup.sh2.2 Adapter 优缺点| 优点 | 缺点 ||------|------|| 无需编写自定义代码 | 灵活性较低 || 开箱即用 | 仅支持常见目标存储 || 配置简单 | 复杂场景可能无法满足需求 || 社区支持良好 | 特定需求可能需要自行扩展 |3. Client 自定义集成方案当 Adapter 无法满足特定需求时可以通过自定义 Client 实现更灵活的数据同步。3.1 Client 集成步骤添加依赖xmldependencygroupIdcom.alibaba.otter/groupIdartifactIdcanal-client/artifactIdversion1.1.5/version/dependency实现客户端代码javaimport 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 CanalClientExample {public static void main(String[] args) {// 创建连接CanalConnector connector CanalConnectors.newSingleConnector(new InetSocketAddress(127.0.0.1, 11111),example,,);try {// 连接connector.connect();// 订阅所有表connector.subscribe(.\\..);// 循环获取数据while (true) {// 获取指定数量的数据Message message connector.getWithoutAck(100);long batchId message.getId();if (batchId -1 || message.getEntries().isEmpty()) {try {Thread.sleep(1000);} catch (InterruptedException e) {e.printStackTrace();}continue;}// 处理数据ListCanalEntry.Entry entries message.getEntries();for (CanalEntry.Entry entry : entries) {if (entry.getEntryType() CanalEntry.EntryType.TRANSACTIONBEGIN|| entry.getEntryType() CanalEntry.EntryType.TRANSACTIONEND) {continue;}CanalEntry.RowChange rowChange null;try {rowChange CanalEntry.RowChange.parseFrom(entry.getStoreValue());} catch (Exception e) {throw new RuntimeException(ERROR ## parser the entry exception!, e);}// 处理每条变更数据for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {// 处理插入数据if (rowChange.getEventType() CanalEntry.EventType.INSERT) {System.out.println(INSERT:);printColumn(rowData.getAfterColumnsList());}// 处理更新数据else if (rowChange.getEventType() CanalEntry.EventType.UPDATE) {System.out.println(UPDATE:);printColumn(rowData.getAfterColumnsList());}// 处理删除数据else if (rowChange.getEventType() CanalEntry.EventType.DELETE) {System.out.println(DELETE:);printColumn(rowData.getBeforeColumnsList());}}}// 确认处理完成connector.ack(batchId);}} finally {connector.disconnect();}}private static void printColumn(ListCanalEntry.Column columns) {for (CanalEntry.Column column : columns) {System.out.println(column.getName() : column.getValue() update column.getUpdated());}}}4. 全量初始化常见问题与解决方案在全量初始化过程中可能会遇到各种问题以下是一些常见问题及解决方案4.1 数据一致性问题问题全量数据快照与后续增量数据之间的时间点不一致导致数据不一致。解决方案使用事务性全量同步确保全量数据的一致性记录全量同步完成的时间点过滤该时间点之前的增量数据使用 CanalExtended 模式确保数据的一致性4.2 大表同步问题问题大表全量同步时间长影响业务性能。解决方案分区表同步按业务分区同步减少单次同步数据量低峰期同步在业务低峰期执行全量同步并行同步多个线程并行同步不同表或分区public class ParallelFullSync { private ExecutorService executor Executors.newFixedThreadPool(5); public void parallelFullSync(ListString tables) { // 使用并行流同步多个表 tables.parallelStream().forEach(table - { try { syncTable(table); } catch (Exception e) { System.err.println(同步表 table 失败: e.getMessage()); } }); } private void syncTable(String table) { // 获取表数据 ListMapString, Object tableData fetchTableData(table); // 写入目标数据库 writeToTarget(table, tableData); } }4.3 数据冲突问题问题主键重复或其他数据冲突问题。解决方案使用 REPLACE INTO 或 ON DUPLICATE KEY UPDATE 处理冲突使用临时表同步完成后替换目标表-- 使用 REPLACE INTO 处理冲突 REPLACE INTO target_table (id, name, age) VALUES (1, Alice, 25); -- 使用临时表同步 CREATE TEMPORARY TABLE temp_table LIKE source_table; INSERT INTO temp_table SELECT * FROM source_table; RENAME TABLE target_table TO target_table_old, temp_table TO target_table; DROP TABLE target_table_old;5. 实战案例与注意事项下面是一个简单的 Canal 客户端实现用于全量增量同步public class FullAndIncrementSync { // 全量同步方法 public void fullSync() { // 从源数据库获取所有表数据 ListTableData allData fetchAllDataFromSource(); // 将数据写入目标数据库 writeToTarget(allData); } // 增量同步方法 public void incrementSync() { // 使用 CanalClientExample 代码监听增量变更 CanalClientExample client new CanalClientExample(); client.start(); } public void start() { // 1. 执行全量同步 fullSync(); // 2. 启动增量同步 incrementSync(); } }5.1 注意事项MySQL 配置确保 MySQL 开启 binlog 功能inilog-binmysql-binbinlog-formatROWserver-id1授予 canal 用户必要的权限sqlCREATE USER canal% IDENTIFIED BY canal;GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON.TO canal%;FLUSH PRIVILEGES;性能优化根据数据量调整批量操作的大小适当增加并行度但避免对目标系统造成过大压力对于大表考虑分批同步或分区间同步错误处理实现重试机制处理临时性错误记录同步日志便于排查问题实现监控告警及时发现问题启动全量同步获取源数据库全量数据清空目标存储批量写入全量数据全量同步完成启动增量同步Canal Server监听MySQL Binlog解析Binlog数据推送变更数据Client接收数据处理增量变更写入目标存储完整实现代码请参考https://github.com/alibaba/calan
返回列表