跳转至主要内容
版本: 5.0

RocketMQ Connect 概述

RocketMQ Connect 是 RocketMQ 数据集成的重要组件,可以高效、可靠地将数据传入和传出 RocketMQ。它是一个独立的、分布式的、可扩展的、且具有容错能力的系统,具备低延迟、高可靠性、高性能、低代码和强扩展性等特点。它可以实现各种异构数据系统的连接、数据管道构建、ETL、CDC 以及数据湖能力。

RocketMQ Connect Overview

Connector 工作原理

RocketMQ Connect 是一个独立的、分布式的、可扩展的、且具有容错能力的系统,主要为 RocketMQ 提供数据流入和流出各种外部系统的能力。用户无需编程,仅需简单的配置即可使用 RocketMQ Connect,例如将数据从 MySQL 同步到 RocketMQ,只需配置账号密码、连接地址以及需要同步的数据库和表名即可。

Connector 使用场景

构建流式数据管道

RocketMQ Connect使用场景

在业务系统中,利用 MySQL 出色的事务支持来处理数据的增删改,使用 ElasticSearch 和 Solr 来实现强大的搜索能力,或者将生成的业务数据同步到数据分析系统和数据湖(如 Hudi)中进行进一步处理,从而使数据产生更高的价值。使用 RocketMQ Connect,可以轻松实现这样的数据管道能力。仅需配置三个任务:第一个任务是从 MySQL 获取数据,第二和第三个任务是从 RocketMQ 消费数据到 ElasticSearch 和 Hudi。通过配置这三个任务,即实现了从 MySQL 到 ElasticSearch 和从 MySQL 到 Hudi 的两条数据管道,这不仅能满足业务中的事务需求,还能满足搜索需求,并能构建数据湖。

CDC

CDC 作为 ETL 模式之一,可以近乎实时地捕获数据库的 INSERT、UPDATE、DELETE 变更。RocketMQ Connect 通过连接器(Connector)进行数据流传输,具有高可用和低延迟的特性,可以轻松实现 CDC。

Connector 部署

创建 Connector 通常是通过配置来完成的。Connector 一般包含逻辑上的 Connector 和执行数据复制的 Task(物理线程),如下图所示,展示了两个 Connector 连接器及其对应的运行 Task 任务。

RocketMQ Connect任务模型1

一个 Connector 也可以同时运行多个任务以增加其并行度。例如,下图中的 Hudi Sink Connector 有 2 个任务,每个任务处理不同的分片数据,从而增加了 Connector 的并行度并提高了处理性能。

RocketMQ Connect任务模型2

RocketMQ Connect Worker 支持集群和单机两种运行模式。在集群模式下,顾名思义,会有多个 Worker 节点,建议至少使用 2 个 Worker 节点以组成高可用集群。集群的配置信息、偏移量(offset)信息和状态信息都存储在指定的 RocketMQ Topic 中。新的 Worker 节点会获取这些配置、偏移量和状态信息,并触发负载均衡,重新分配集群中的任务以达到平衡状态;当减少 Worker 节点或 Worker 节点宕机时,也会触发负载均衡,以确保集群中的所有任务都能在存活的节点上正常运行。

RocketMQ Connect部署模型集群

在单机模式下,Connector 任务在单机上运行,Worker 本身不具备高可用性,任务的偏移量信息保存在本地。它适用于没有高可用需求或不需要 Worker 自身保证高可用的场景,例如部署在 K8s 集群中,此时由 K8s 集群来提供保障。

RocketMQ Connect部署模型单机