欢迎光临
我们一直在努力

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

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 架构流程图

    #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 最佳实践

  • 位点和容错:确保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 配置:
    • 根据数据量和查询模式选择合适的表结构和分区策略
    • 调整批处理大小,平衡实时性和性能
    • 设置合适的压缩和编码策略,优化存储空间
  • 数据一致性保障:
    • 实现同步数据的校验机制
    • 定期进行数据一致性检查
    • 设置合理的重试和回滚机制
  • 监控与运维:
    • 建立完善的监控告警机制
    • 定期查看同步延迟情况
    • 准备应急方案,应对可能的故障情况
    赞(0)
    未经允许不得转载:171主机测评 » Canal 与 HBase/Hudi 实时入湖:MySQL 数据实时同步到数据湖的架构实践
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址