news 2026/9/4 13:25:52

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Canal 与 HBase/Hudi 实时入湖:MySQL 数据实时同步到数据湖的架构实践

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

1. Canal 简介

Canal是阿里巴巴开源的一款基于数据库增量日志解析的中间件,主要用于解决数据库实时同步问题。它通过解析MySQL的binlog日志,实现对数据库变更的捕获和传递,为数据实时同步提供了高效可靠的技术基础。

1.1 Canal 工作原理

Canal的工作原理主要包括以下几个步骤:

  1. 连接MySQL:Canal作为MySQL的从库,连接到MySQL主库并请求binlog日志
  2. 解析binlog:解析MySQL的binlog,获取数据库变更事件
  3. 数据转换:将binlog事件转换为Canal定义的消息格式
  4. 消息发送:将变更消息发送给下游消费者

1.2 Canal 核心组件

  • server:Canal的核心服务,负责连接MySQL并解析binlog
  • instance:Canal的实例,每个实例对应一个数据源的同步任务
  • meta manager:管理同步位点信息,确保数据不丢失
  • sink:数据消费模块,负责将变更数据发送到目标系统

2. Hudi/HBase 实时入湖架构

将MySQL数据同步到数据湖,通常采用HBase或Hudi作为中间存储,实现湖仓一体的架构设计。下面详细介绍这两种方案的架构设计。

2.1 基于 HBase 的实时入湖架构

基于HBase的实时入湖架构主要包括以下组件:

  1. MySQL:源数据存储,开启binlog功能
  2. Canal Server:捕获MySQL的binlog日志
  3. HBase:作为中间存储层,提供快速读写能力
  4. 数据湖:最终数据存储,如HDFS、S3等
  5. 数据处理应用:负责将HBase中的数据同步到数据湖

该架构的特点是利用HBase的随机读写能力,为数据湖提供实时查询能力,同时保证数据一致性。

2.2 基于 Hudi 的实时入湖架构

基于Hudi的实时入湖架构是更为现代的湖仓一体化方案:

  1. MySQL:源数据存储,开启binlog功能
  2. Canal Server:捕获MySQL的binlog日志
  3. Kafka:消息队列,缓存变更数据
  4. Hudi:提供数据湖上的ACID事务和增量处理能力
  5. 数据湖:基于HDFS、S3等存储系统的数据湖

该架构的特点是Hudi直接在数据湖上提供类似数据库的事务能力,实现了存储计算分离和湖仓一体的架构。

2.3 架构对比

| 特性 | HBase 架构 | Hudi 架构 |

|------|------------|------------|

| 数据一致性 | 强一致性 | 最终一致性 |

| 实时性 | 高 | 高 |

| 查询能力 | 强 | 中等 |

| 存储成本 | 高 | 低 |

| 扩展性 | 中等 | 高 |

| 适用场景 | 需要强一致性查询 | 需要低存储成本和高扩展性 |

3. 实施步骤与实践经验

基于上述架构,以下是具体的实施步骤和实践经验。

3.1 Canal 部署与配置

  1. 安装 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
  1. 配置 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 配置
  1. 创建 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());
  1. 编写消费逻辑
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 配置
  1. 创建 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);
  1. 编写消费逻辑
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 架构流程图

binlog日志增量数据变更实时数据流批量数据同步数据查询与分析实时查询

MySQL 数据库

Canal Server

消息队列 Kafka

HBase/Hudi

数据湖

BI/报表系统

实时查询应用

3.4 最佳实践

  1. 位点和容错:确保Canal的正确记录位点,避免数据丢失
  2. 监控告警:建立完善的监控机制,及时发现同步异常
  3. 性能调优:根据业务场景调整Canal、HBase/Hudi的参数配置
  4. 数据一致性:确保同步过程中数据的一致性,特别是关键业务数据
  5. 数据质量:建立数据质量校验机制,确保同步数据的正确性

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 注意事项

  1. MySQL 配置
  • 确保 MySQL 开启 binlog 模式
  • 设置 binlog_format 为 ROW 格式
  • 合理设置 binlog 相关参数,避免磁盘空间不足
  1. Canal 配置
  • 根据实际情况调整内存和线程参数
  • 设置合理的位点信息保存策略
  • 配置合适的过滤规则,避免同步过多无用数据
  1. HBase/Hudi 配置
  • 根据数据量和查询模式选择合适的表结构和分区策略
  • 调整批处理大小,平衡实时性和性能
  • 设置合适的压缩和编码策略,优化存储空间
  1. 数据一致性保障
  • 实现同步数据的校验机制
  • 定期进行数据一致性检查
  • 设置合理的重试和回滚机制
  1. 监控与运维
  • 建立完善的监控告警机制
  • 定期查看同步延迟情况
  • 准备应急方案,应对可能的故障情况
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/4 13:24:11

当国风品牌遇见商业插画: 让创意落地真实的品牌商业场景

核心摘要 • 画星人与茶颜悦色的合作&#xff0c;本质上是一次创意人才与真实品牌需求之间的连接。 • 合作不只是“提供插画”&#xff0c;更重要的是把创作者带入真实商业项目&#xff0c;让插画能力参与品牌视觉与内容创新。 • 从需求理解、创意产出到项目落地&#xff0c…

作者头像 李华
网站建设 2026/9/4 13:23:48

LLM代码审查也会“表扬”Bug?一场静默失败评测实验与优化指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/4 13:23:42

算力崛起,真正底座其实是电网

《算力崛起&#xff0c;真正底座其实是电网》——未来比的不是谁发电更多&#xff0c;而是谁能把绿电接住、送稳、用准过去十年&#xff0c;算力是AI的杠杆&#xff1b;未来十年&#xff0c;电网可能才是算力的天花板。“十五五”期间&#xff0c;我国电力需求预计仍将年均增长…

作者头像 李华
网站建设 2026/9/4 13:19:22

高速PCB设计全流程解析:叠层、阻抗与信号完整性分析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华