ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Canal 与 HBase/Hudi 实时入湖:MySQL 数据实时同步到数据湖的架构实践

Canal 与 HBase/Hudi 实时入湖:MySQL 数据实时同步到数据湖的架构实践 Canal 与 HBase/Hudi 实时入湖MySQL 数据实时同步到数据湖的架构实践1. Canal 简介Canal是阿里巴巴开源的一款基于数据库增量日志解析的中间件主要用于解决数据库实时同步问题。它通过解析MySQL的binlog日志实现对数据库变更的捕获和传递为数据实时同步提供了高效可靠的技术基础。1.1 Canal 工作原理Canal的工作原理主要包括以下几个步骤连接MySQLCanal作为MySQL的从库连接到MySQL主库并请求binlog日志解析binlog解析MySQL的binlog获取数据库变更事件数据转换将binlog事件转换为Canal定义的消息格式消息发送将变更消息发送给下游消费者1.2 Canal 核心组件serverCanal的核心服务负责连接MySQL并解析binloginstanceCanal的实例每个实例对应一个数据源的同步任务meta manager管理同步位点信息确保数据不丢失sink数据消费模块负责将变更数据发送到目标系统2. Hudi/HBase 实时入湖架构将MySQL数据同步到数据湖通常采用HBase或Hudi作为中间存储实现湖仓一体的架构设计。下面详细介绍这两种方案的架构设计。2.1 基于 HBase 的实时入湖架构基于HBase的实时入湖架构主要包括以下组件MySQL源数据存储开启binlog功能Canal Server捕获MySQL的binlog日志HBase作为中间存储层提供快速读写能力数据湖最终数据存储如HDFS、S3等数据处理应用负责将HBase中的数据同步到数据湖该架构的特点是利用HBase的随机读写能力为数据湖提供实时查询能力同时保证数据一致性。2.2 基于 Hudi 的实时入湖架构基于Hudi的实时入湖架构是更为现代的湖仓一体化方案MySQL源数据存储开启binlog功能Canal Server捕获MySQL的binlog日志Kafka消息队列缓存变更数据Hudi提供数据湖上的ACID事务和增量处理能力数据湖基于HDFS、S3等存储系统的数据湖该架构的特点是Hudi直接在数据湖上提供类似数据库的事务能力实现了存储计算分离和湖仓一体的架构。2.3 架构对比| 特性 | HBase 架构 | Hudi 架构 ||------|------------|------------|| 数据一致性 | 强一致性 | 最终一致性 || 实时性 | 高 | 高 || 查询能力 | 强 | 中等 || 存储成本 | 高 | 低 || 扩展性 | 中等 | 高 || 适用场景 | 需要强一致性查询 | 需要低存储成本和高扩展性 |3. 实施步骤与实践经验基于上述架构以下是具体的实施步骤和实践经验。3.1 Canal 部署与配置安装 Canal# 下载 Canal wget https://github.com/alibaba/canal/releases/download/canal-1.1.4/canal.deployer-1.1.4.tar.gz tar -zxvf canal.deployer-1.1.4.tar.gz cd canal.deployer # 修改配置文件 vim conf/example/instance.properties # 启动 Canal bin/startup.sh配置 MySQL-- 创建用户 CREATE USER canal% IDENTIFIED BY canal; -- 授权 GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; -- 刷新权限 FLUSH PRIVILEGES;3.2 HBase/Hudi 配置3.2.1 HBase 配置创建 HBase 表// 创建 HBase 表 Connection connection ConnectionFactory.createConnection(admin.getConfiguration()); Table table connection.getTable(TableName.valueOf(user_table)); // 定义表结构 TableDescriptorBuilder tableDescriptorBuilder TableDescriptorBuilder.newBuilder(TableName.valueOf(user_table)); ColumnFamilyDescriptorBuilder columnFamilyDescriptorBuilder ColumnFamilyDescriptorBuilder.newBuilder(Bytes.toBytes(info)); tableDescriptorBuilder.setColumnFamily(columnFamilyDescriptorBuilder.build()); // 创建表 admin.createTable(tableDescriptorBuilder.build());编写消费逻辑public class HBaseConsumer implements CanalEventSinkCanalEntry.Entry { Override public void sink(ListCanalEntry.Entry entries, Context context) { Connection connection null; try { connection ConnectionFactory.createConnection(); Table table connection.getTable(TableName.valueOf(user_table)); for (CanalEntry.Entry entry : entries) { if (entry.getEntryType() CanalEntry.EntryType.ROWDATA) { CanalEntry.RowChange rowChange CanalEntry.RowChange.parseFrom(entry.getStoreValue()); for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { if (rowChange.getEventType() CanalEntry.EventType.INSERT) { // 处理插入操作 Put put convertToPut(rowData.getAfterColumnsList()); table.put(put); } else if (rowChange.getEventType() CanalEntry.EventType.UPDATE) { // 处理更新操作 Put put convertToPut(rowData.getAfterColumnsList()); table.put(put); } else if (rowChange.getEventType() CanalEntry.EventType.DELETE) { // 处理删除操作 Delete delete convertToDelete(rowData.getBeforeColumnsList()); table.delete(delete); } } } } table.flush(); } catch (Exception e) { throw new RuntimeException(Error while processing Canal event, e); } finally { if (connection ! null) { try { connection.close(); } catch (IOException e) { // 忽略关闭异常 } } } } private Put convertToPut(ListCanalEntry.Column columns) { // 将 Canal 列转换为 HBase Put } private Delete convertToDelete(ListCanalEntry.Column columns) { // 将 Canal 列转换为 HBase Delete } }3.2.2 Hudi 配置创建 Hudi 表// 创建 Hudi 配置 MapString, String configs new HashMap(); configs.put(hoodie.table.payload.class, org.apache.hoodie.client.transaction.lock.ZookeeperBasedLockProvider); configs.put(hoodie.table.name, user_table); configs.put(hoodie.table.type, COPY_ON_WRITE); configs.put(hoodie.table.payload.class, org.apache.hoodie.common.model.PartialUpdateAvroPayload); configs.put(hoodie.cleaner.commits.retained, 10); configs.put(hoodie.timeline.server.port, 10000); // 创建 Hudi 表 HoodieWriteConfig writeConfig HoodieWriteConfig.newBuilder() .withPath(s3://your-bucket/path/to/table) .withSchema(id:int,name:string,age:int) .withWriteConcurrencyMode(HoodieWriteConcurrencyMode.OPTIMISTIC_CONCURRENCY_CONTROL) .withBulkInsertSortMemoryInBytes(1024 * 1024 * 128) .withBulkInsertSortShuffleInput(1024 * 1024 * 128) .withBulkInsertSortMemory(1024 * 1024 * 128) .build(); HoodieTable table HoodieTable.create(configs, writeConfig);编写消费逻辑public class HudiConsumer implements CanalEventSinkCanalEntry.Entry { Override public void sink(ListCanalEntry.Entry entries, Context context) { HoodieWriteConfig writeConfig // 初始化配置 JavaSparkSession spark JavaSparkSession.builder().appName(HudiCanalSync).getOrCreate(); for (CanalEntry.Entry entry : entries) { if (entry.getEntryType() CanalEntry.EntryType.ROWDATA) { CanalEntry.RowChange rowChange CanalEntry.RowChange.parseFrom(entry.getStoreValue()); for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { if (rowChange.getEventType() CanalEntry.EventType.INSERT) { // 处理插入操作 DatasetRow df convertToDataFrame(rowData.getAfterColumnsList(), spark); HoodieWriteResult result df.write() .format(org.apache.hudi) .options(writeConfig.getProps()) .option(hoodie.table.name, user_table) .mode(Append) .save(); } else if (rowChange.getEventType() CanalEntry.EventType.UPDATE) { // 处理更新操作 DatasetRow df convertToDataFrame(rowData.getAfterColumnsList(), spark); HoodieWriteResult result df.write() .format(org.apache.hudi) .options(writeConfig.getProps()) .option(hoodie.table.name, user_table) .mode(Append) .save(); } else if (rowChange.getEventType() CanalEntry.EventType.DELETE) { // 处理删除操作 DatasetRow df convertToDataFrame(rowData.getBeforeColumnsList(), spark); HoodieWriteResult result df.write() .format(org.apache.hudi) .options(writeConfig.getProps()) .option(hoodie.table.name, user_table) .mode(delete) .save(); } } } } spark.stop(); } private DatasetRow convertToDataFrame(ListCanalEntry.Column columns, JavaSparkSession spark) { // 将 Canal 列转换为 Spark DataFrame } }3.3 架构流程图binlog日志增量数据变更实时数据流批量数据同步数据查询与分析实时查询MySQL 数据库Canal Server消息队列 KafkaHBase/Hudi数据湖BI/报表系统实时查询应用3.4 最佳实践位点和容错确保Canal的正确记录位点避免数据丢失监控告警建立完善的监控机制及时发现同步异常性能调优根据业务场景调整Canal、HBase/Hudi的参数配置数据一致性确保同步过程中数据的一致性特别是关键业务数据数据质量建立数据质量校验机制确保同步数据的正确性4. 最小示例与注意事项4.1 最小示例以下是一个基于CanalHudi的简单同步示例public class CanalHudiExample { public static void main(String[] args) { // 创建Canal客户端 CanalConnector connector CanalConnectors.newSingleConnector( new InetSocketAddress(127.0.0.1, 11111), example, canal, canal); // 创建Hudi消费者 HudiConsumer hudiConsumer new HudiConsumer(); try { connector.connect(); connector.subscribe(.*\\..*); // 订阅所有库的所有表 connector.rollback(100L); // 回滚到未确认位置 while (true) { Message message connector.getWithoutAck(100); long batchId message.getId(); if (batchId -1 || message.getEntries().isEmpty()) { Thread.sleep(1000); continue; } // 处理消息 ListCanalEntry.Entry entries message.getEntries(); hudiConsumer.sink(entries, null); // 提交确认 connector.ack(batchId); } } catch (Exception e) { e.printStackTrace(); } finally { connector.disconnect(); } } }4.2 注意事项MySQL 配置确保 MySQL 开启 binlog 模式设置 binlog_format 为 ROW 格式合理设置 binlog 相关参数避免磁盘空间不足Canal 配置根据实际情况调整内存和线程参数设置合理的位点信息保存策略配置合适的过滤规则避免同步过多无用数据HBase/Hudi 配置根据数据量和查询模式选择合适的表结构和分区策略调整批处理大小平衡实时性和性能设置合适的压缩和编码策略优化存储空间数据一致性保障实现同步数据的校验机制定期进行数据一致性检查设置合理的重试和回滚机制监控与运维建立完善的监控告警机制定期查看同步延迟情况准备应急方案应对可能的故障情况
返回列表