Canal 与 HBase/Hudi 实时入湖:MySQL 数据实时同步到数据湖的架构实践
1. Canal 简介
Canal是阿里巴巴开源的一款基于数据库增量日志解析的中间件,主要用于解决数据库实时同步问题。它通过解析MySQL的binlog日志,实现对数据库变更的捕获和传递,为数据实时同步提供了高效可靠的技术基础。
1.1 Canal 工作原理
Canal的工作原理主要包括以下几个步骤:
- 连接MySQL:Canal作为MySQL的从库,连接到MySQL主库并请求binlog日志
- 解析binlog:解析MySQL的binlog,获取数据库变更事件
- 数据转换:将binlog事件转换为Canal定义的消息格式
- 消息发送:将变更消息发送给下游消费者
1.2 Canal 核心组件
- server:Canal的核心服务,负责连接MySQL并解析binlog
- instance:Canal的实例,每个实例对应一个数据源的同步任务
- 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 CanalEventSink<CanalEntry.Entry> { @Override public void sink(List<CanalEntry.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(List<CanalEntry.Column> columns) { // 将 Canal 列转换为 HBase Put } private Delete convertToDelete(List<CanalEntry.Column> columns) { // 将 Canal 列转换为 HBase Delete } }3.2.2 Hudi 配置
- 创建 Hudi 表:
// 创建 Hudi 配置 Map<String, 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 CanalEventSink<CanalEntry.Entry> { @Override public void sink(List<CanalEntry.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) { // 处理插入操作 Dataset<Row> 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) { // 处理更新操作 Dataset<Row> 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) { // 处理删除操作 Dataset<Row> 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 Dataset<Row> convertToDataFrame(List<CanalEntry.Column> columns, JavaSparkSession spark) { // 将 Canal 列转换为 Spark DataFrame } }3.3 架构流程图
3.4 最佳实践
- 位点和容错:确保Canal的正确记录位点,避免数据丢失
- 监控告警:建立完善的监控机制,及时发现同步异常
- 性能调优:根据业务场景调整Canal、HBase/Hudi的参数配置
- 数据一致性:确保同步过程中数据的一致性,特别是关键业务数据
- 数据质量:建立数据质量校验机制,确保同步数据的正确性
4. 最小示例与注意事项
4.1 最小示例
以下是一个基于Canal+Hudi的简单同步示例:
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; } // 处理消息 List<CanalEntry.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 配置:
- 根据数据量和查询模式选择合适的表结构和分区策略
- 调整批处理大小,平衡实时性和性能
- 设置合适的压缩和编码策略,优化存储空间
- 数据一致性保障:
- 实现同步数据的校验机制
- 定期进行数据一致性检查
- 设置合理的重试和回滚机制
- 监控与运维:
- 建立完善的监控告警机制
- 定期查看同步延迟情况
- 准备应急方案,应对可能的故障情况