
1 简介RocketMQ Connect是RocketMQ数据集成重要组件可将各种系统中的数据通过高效可靠流的方式流入流出到 RocketMQ它是独立于 RocketMQ 的一个单独的分布式可扩展可容错系统它具备低延时高靠性高性能低代码扩展性强等特点可以实现各种异构数据系统的连接构建数据管道ETLCDC数据湖等能力。RocketMQ EventBridge架起高可伸缩高吞吐的事件通道。 EventBridge可以看作是runtime(connect)的管理台构建逻辑(based domain模块)云事件通道调用runtime api(adapter)创建物理(可运行)事件通道本系列文章分析RocketMQ Connect和RocketMQ EventBridge源码原理为两组件的改造和应用开发提供支持本文是第一部分分析架构服务和组件第二部分分析workerworker connectorworker source taskworker sink task2 关键词CloudEventOpenMessaging OpenMessaging 是 Linux 基金会下一个开源组织致力于制定消息领域的标准3 参考资料4 发布计划M1 事件通道构建控制台Ø 拖拽方式构建事件通道包括source connector/sink connector及其transform链集成rule based transformØ 可视化worker/worker集群状态和启停connector/task启停M2 分布式重构Ø 引入elastic-platform(based zookeeper)实现worker集群管理配置服务集群成员变更分片M3 动态资源管理Ø 引入资源管理器支持k8s资源申请5 Connect本章介绍connect平台5.1 集群架构若干worker组成集群同一集群的worker共同运行同一个事件通道worker间互为负载均衡和故障转移。5.2 逻辑架构上图是事件通道的展开事件通道是逻辑定义每个worker自身运行事件通道的部分任务事件通道多个源和多个目标源与目标之间转换链其中filter转换决定是否流到对应的sink5.3 组件和场景*上面包视图根据个人理解有部分改动startupworker启动支持standalone和分布式两种模式restworker的rest接口输出worker内connector和任务的操作服务集群服务状态服务配置服务位点服务分片服务等workerconnect执行单位封装连接器/任务转换执行状态维护controller控制器起着门面服务的作用持有服务实例plugin类加载器载入jar包类型utilsServiceThread和datasync组件还有些组件本系列没有分析connectortransform5.4 原理分析原理源码分析通过场景用例分析完成包括启动服务worker(包括连接器和任务)系统组件5.4.1 启动(startup)5.4.1.1 worker启动worker启动有两种模式分布式和standalone流程差不多下面分析一下分布式模式启动启动从startup类发起该类从命令行获取配置实例化服务类插件类然后构建控制器控制器实例和初始化各服务和功能类控制器两个职责1. 作为门面服务支持rest服务2. 持有服务类传递给Worker5.4.1.2 worker集群启动目前connect并没有worker集群的启动实现worker节点启动加入集群集群的其他worker收到消费组变更通知触发重分片5.4.2 服务下面详细分析每个服务5.4.2.1 分片服务分片服务负责分派连接器和任务到workerRelanceService分片服务继承ServiceThread获得定时执行和唤醒执行的特性RebanceImpl/AllocateConnAndTaskStrategy分片的实现类支持不同的分片策略分片完成连接器/任务分派到位调用Worker重新启动连接器和任务WorkerStatusListenerImpl实现WorkerStatusListener接口监听集群变化唤醒分片服务驱动重新分片ConnectorConnectorConfigChangeListenerImpl实现ConnectorConfigUpdateListener接口监听connector/task配置变更唤醒分片服务驱动重新分片5.4.2.2 分片策略分片策略有2点关键1. 分片服务集成ServiceThread支持唤醒执行和定时执行而且定时时间间隔1秒原因可能是集群变更事件有丢失或者执行失败可能引起任务丢失频繁的分片要求分片策略有很高的稳定性尽量减少连接器和任务的转移2. 分片每个worker进行分派属于自身worker的部分给自己因此分片策略需与worker相关不遗漏不重复的分派connect自带两个实现默认实现workerId排序然后哈希分派另一个是哈希一致性这里不详细分析5.4.2.3 集群服务集群服务实现比较简单职责是监听集群成员变化对外输出WorkerStatusListener监听接口集群成员变更事件驱动执行分片服务监听集群成员变化是通过订阅消息引擎的消费组成员变更事件规划引入zookeeper重新实现集群服务配置服务5.4.2.4 配置服务(config)配置服务负责连接器/任务配置存储同步配置新增变更连接器/任务启停ConfigManagementService配置管理服务接口有localmemoryrocketmq实现适配standalone和分布式模式memory只能在standalonelocal和rocketmq可用作分布式模式*DataSync组件负责分发请求到集群的所有worker请求分两类配置变更connector/task启停操作两类请求有重叠配置变更最终通过重新分片反映到connector/task对于配置变更worker不直接处理通过DataSync分发变更到集群当然包括自身在消息消费中处理这样所有配置在每个worker有完整备份用于故障转移新worker同步对于启停操作connector/task分派到不同workerworker不直接处理通过DataSync分发变更到集群各worker收到后识别是否自身处理ConnectorConfigUpdateListener配置变更通知接口目前有一个实现该实现分配服务介绍过唤醒分片服务驱动重新分片Ø 主要业务DataSync组件分析可知主要业务集中在消息消费方法TARGET_STATE_PREFIX/CONNECTOR_PREFIX/TASK_PREFIX/DELETE_CONNECTOR_PREFIX系统通过消息key前缀识别分流处理具体业务不详细分析local实现起始时发送START_SIGNAL收集其他worker的配置数据5.4.2.5 位点服务(position)位点服务记录source/sink处理位置用于分页和容错“套路”与配置服务一样原理不重复分析PositionUpdateListener 位点变更监听目前没有使用Ø 主要业务分析ONLINE 是新节点起来请求同步集群的其他节点数据获得全量数据5.4.2.6 状态服务(state)状态服务跟踪connector/task的状态“套路”与上面两个服务一样原理不重复分析Ø 状态模型WrapperStatusListener该接口是聚合连接器和任务监听实现连接器和任务通过该接口报告给状态服务ConnAndTaskStatus状态服务实现都拥有KeyValueStore但并没有使用而是使用ConnAndTaskStatus保存连接器和任务的状态Ø 主业务分析5.4.3 worker第二部分分析5.4.4 资源管理目前connect没有资源管理计算资源即worker节点worker节点相同消费组是同一个集群作为计算资源分派任务规划引入资源管理器对接k8s动态申请资源5.4.5 功能组件5.4.5.1 rest组件worker内置rest服务RestHandle输出worker连接器/任务新建配置和变更等服务实现依赖控制器5.4.5.2 服务唤醒和等待组件(ServiceThread)服务继承ServiceThread获得定时和唤醒执行服务逻辑的特性5.4.5.3 同步组件datasync底层以消息引擎实现用于数据同步和命令分发例如任务配置变更通知集群的其他worker其他worker保存到本地每份数据都得到复制用于任务故障转移时恢复5.4.5.4 store组件上一章提到位点状态配置的存储store组件负责该功能组件利用泛型和序列化器支持不同的数据的存储实现只有基于内存和文件均为本地存储5.4.5.5 统计(stat)组件TBD6 EventBridge运行架构总体来说EventBridge可以看作是runtime(connect)的管理台构建基于domain的逻辑事件通道远程调用runtime api创建物理(可运行)事件通道组件架构core 逻辑事件通道的模型rpc 对应worker的RestHandle构建/启停connector等操作web 逻辑事件通道模型的构建基于spring boot的web应用其中webrpc持久属于adapter子系统不同的runtim有对应的实现domain是抽象的事件通道模型集成弹性资源随着弹性资源平台的完成集成到弹性资源提供高弹性伸缩https://blog.csdn.net/szlhj/category_12446958.html