Canal 与 HBase/Hudi 实时入湖:MySQL 数据实时同步到数据湖的架构实践
1. Canal 简介
Canal是阿里巴巴开源的一款基于数据库增量日志解析的中间件,主要用于解决数据库实时同步问题。它通过解析MySQL的binlog日志,实现对数据库变更的捕获和传递,为数据实时同步提供了高效可靠的技术基础。
1.1 Canal 工作原理
Canal的工作原理主要包括以下几个步骤:
1.2 Canal 核心组件
- server:Canal的核心服务,负责连接MySQL并解析binlog
- instance:Canal的实例,每个实例对应一个数据源的同步任务
- meta manager:管理同步位点信息,确保数据不丢失
- sink:数据消费模块,负责将变更数据发送到目标系统
2. Hudi/HBase 实时入湖架构
将MySQL数据同步到数据湖,通常采用HBase或Hudi作为中间存储,实现湖仓一体的架构设计。下面详细介绍这两种方案的架构设计。
2.1 基于 HBase 的实时入湖架构
基于HBase的实时入湖架构主要包括以下组件:
该架构的特点是利用HBase的随机读写能力,为数据湖提供实时查询能力,同时保证数据一致性。
2.2 基于 Hudi 的实时入湖架构
基于Hudi的实时入湖架构是更为现代的湖仓一体化方案:
该架构的特点是Hudi直接在数据湖上提供类似数据库的事务能力,实现了存储计算分离和湖仓一体的架构。
2.3 架构对比
| 特性 | HBase 架构 | Hudi 架构 |
|——|————|————|
| 数据一致性 | 强一致性 | 最终一致性 |
| 实时性 | 高 | 高 |
| 查询能力 | 强 | 中等 |
| 存储成本 | 高 | 低 |
| 扩展性 | 中等 | 高 |
| 适用场景 | 需要强一致性查询 | 需要低存储成本和高扩展性 |
3. 实施步骤与实践经验
基于上述架构,以下是具体的实施步骤和实践经验。
3.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
— 创建用户
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 表
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 配置
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 架构流程图
#publish-mermaid-1788489137752-0{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#publish-mermaid-1788489137752-0 .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#publish-mermaid-1788489137752-0 .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#publish-mermaid-1788489137752-0 .error-icon{fill:#552222;}#publish-mermaid-1788489137752-0 .error-text{fill:#552222;stroke:#552222;}#publish-mermaid-1788489137752-0 .edge-thickness-normal{stroke-width:1px;}#publish-mermaid-1788489137752-0 .edge-thickness-thick{stroke-width:3.5px;}#publish-mermaid-1788489137752-0 .edge-pattern-solid{stroke-dasharray:0;}#publish-mermaid-1788489137752-0 .edge-thickness-invisible{stroke-width:0;fill:none;}#publish-mermaid-1788489137752-0 .edge-pattern-dashed{stroke-dasharray:3;}#publish-mermaid-1788489137752-0 .edge-pattern-dotted{stroke-dasharray:2;}#publish-mermaid-1788489137752-0 .marker{fill:#333333;stroke:#333333;}#publish-mermaid-1788489137752-0 .marker.cross{stroke:#333333;}#publish-mermaid-1788489137752-0 svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#publish-mermaid-1788489137752-0 p{margin:0;}#publish-mermaid-1788489137752-0 .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#publish-mermaid-1788489137752-0 .cluster-label text{fill:#333;}#publish-mermaid-1788489137752-0 .cluster-label span{color:#333;}#publish-mermaid-1788489137752-0 .cluster-label span p{background-color:transparent;}#publish-mermaid-1788489137752-0 .label text,#publish-mermaid-1788489137752-0 span{fill:#333;color:#333;}#publish-mermaid-1788489137752-0 .node rect,#publish-mermaid-1788489137752-0 .node circle,#publish-mermaid-1788489137752-0 .node ellipse,#publish-mermaid-1788489137752-0 .node polygon,#publish-mermaid-1788489137752-0 .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#publish-mermaid-1788489137752-0 .rough-node .label text,#publish-mermaid-1788489137752-0 .node .label text,#publish-mermaid-1788489137752-0 .image-shape .label,#publish-mermaid-1788489137752-0 .icon-shape .label{text-anchor:middle;}#publish-mermaid-1788489137752-0 .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#publish-mermaid-1788489137752-0 .rough-node .label,#publish-mermaid-1788489137752-0 .node .label,#publish-mermaid-1788489137752-0 .image-shape .label,#publish-mermaid-1788489137752-0 .icon-shape .label{text-align:center;}#publish-mermaid-1788489137752-0 .node.clickable{cursor:pointer;}#publish-mermaid-1788489137752-0 .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#publish-mermaid-1788489137752-0 .arrowheadPath{fill:#333333;}#publish-mermaid-1788489137752-0 .edgePath .path{stroke:#333333;stroke-width:1px;}#publish-mermaid-1788489137752-0 .flowchart-link{stroke:#333333;fill:none;}#publish-mermaid-1788489137752-0 .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#publish-mermaid-1788489137752-0 .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#publish-mermaid-1788489137752-0 .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#publish-mermaid-1788489137752-0 .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#publish-mermaid-1788489137752-0 .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#publish-mermaid-1788489137752-0 .cluster text{fill:#333;}#publish-mermaid-1788489137752-0 .cluster span{color:#333;}#publish-mermaid-1788489137752-0 div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#publish-mermaid-1788489137752-0 .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#publish-mermaid-1788489137752-0 rect.text{fill:none;stroke-width:0;}#publish-mermaid-1788489137752-0 .icon-shape,#publish-mermaid-1788489137752-0 .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#publish-mermaid-1788489137752-0 .icon-shape p,#publish-mermaid-1788489137752-0 .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#publish-mermaid-1788489137752-0 .icon-shape .label rect,#publish-mermaid-1788489137752-0 .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#publish-mermaid-1788489137752-0 .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#publish-mermaid-1788489137752-0 .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#publish-mermaid-1788489137752-0 .node .neo-node{stroke:#9370DB;}#publish-mermaid-1788489137752-0 [data-look=\”neo\”].node rect,#publish-mermaid-1788489137752-0 [data-look=\”neo\”].cluster rect,#publish-mermaid-1788489137752-0 [data-look=\”neo\”].node polygon{stroke:#9370DB;filter:drop-shadow(1px 2px 2px rgba(185, 185, 185, 1));}#publish-mermaid-1788489137752-0 [data-look=\”neo\”].swimlane.cluster rect{filter:none;}#publish-mermaid-1788489137752-0 [data-look=\”neo\”].node path{stroke:#9370DB;stroke-width:1px;}#publish-mermaid-1788489137752-0 [data-look=\”neo\”].node .outer-path{filter:drop-shadow(1px 2px 2px rgba(185, 185, 185, 1));}#publish-mermaid-1788489137752-0 [data-look=\”neo\”].node .neo-line path{stroke:#9370DB;filter:none;}#publish-mermaid-1788489137752-0 [data-look=\”neo\”].node circle{stroke:#9370DB;filter:drop-shadow(1px 2px 2px rgba(185, 185, 185, 1));}#publish-mermaid-1788489137752-0 [data-look=\”neo\”].node circle .state-start{fill:#000000;}#publish-mermaid-1788489137752-0 [data-look=\”neo\”].icon-shape .icon{fill:#9370DB;filter:drop-shadow(1px 2px 2px rgba(185, 185, 185, 1));}#publish-mermaid-1788489137752-0 [data-look=\”neo\”].icon-shape .icon-neo path{stroke:#9370DB;filter:drop-shadow(1px 2px 2px rgba(185, 185, 185, 1));}#publish-mermaid-1788489137752-0 :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}binlog 日志增量数据变更实时数据流批量数据同步数据查询与分析实时查询
MySQL 数据库
Canal Server
消息队列 Kafka
HBase/Hudi
数据湖
BI/报表系统
实时查询应用
3.4 最佳实践
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 开启 binlog 模式
- 设置 binlog_format 为 ROW 格式
- 合理设置 binlog 相关参数,避免磁盘空间不足
- 根据实际情况调整内存和线程参数
- 设置合理的位点信息保存策略
- 配置合适的过滤规则,避免同步过多无用数据
- 根据数据量和查询模式选择合适的表结构和分区策略
- 调整批处理大小,平衡实时性和性能
- 设置合适的压缩和编码策略,优化存储空间
- 实现同步数据的校验机制
- 定期进行数据一致性检查
- 设置合理的重试和回滚机制
- 建立完善的监控告警机制
- 定期查看同步延迟情况
- 准备应急方案,应对可能的故障情况

