跳转至主要内容
版本: 5.0

RocketMQ Streams 概述

RocketMQ Streams 是一个基于 RocketMQ 的轻量级流计算引擎。它以 SDK 依赖的形式应用,无需部署复杂的流计算服务器,具有资源高效、易于扩展以及流计算算子丰富的特点。

架构

总体架构

数据由 RocketMQ-streams 从 RocketMQ 中消费、处理,并最终写回 RocketMQ。

总体架构

数据由 RocketMQ Consumer 消费,进入处理拓扑中被算子处理。如果流处理任务包含 keyBy 算子,数据需要根据 Key 进行分组并写入 shuffle topic。后续算子从 shuffle topic 中消费。如果还存在诸如 count 等有状态算子,计算过程需要读写状态 topic。计算完成后,结果被写回 RocketMQ。

消费模型

img_2.png

计算实例实际依赖 Rocket-streams SDK 的客户端。因此,计算实例消费 MQ,依赖 RocketMQ 的重平衡(rebalance)分配机制。计算实例的总数不能超过消费 MQ 的总数,否则部分计算实例将处于等待状态,无法消费数据。

一个计算实例可以消费多个 MQ,而在一个实例内部,仅存在一个计算拓扑图。

状态

img_3.png

对于有状态算子(如 count),必须先进行分组(grouping)再求和。分组算子 keyBy 会根据分组键将数据重新写入 RocketMQ,并确保具有相同键的数据被写入同一个分区(此过程称为 shuffle),以确保相同键的数据被同一个消费者消费。状态通过 RocksDB 进行本地加速,并通过 RocketMQ 进行远程持久化。

扩缩容能力

img.png

当计算实例从 3 个减少到 2 个时,借助 RocketMQ 集群消费模式下的重平衡功能,消费的 MQ 将在计算实例之间重新分配。原本由实例 1 消费的 MQ2 和 MQ3 被分配给实例 2 和实例 3,这两个 MQ 的状态数据也需要迁移到实例 2 和实例 3。这也意味着状态数据是按照原始数据分区 MQ 进行保存的;扩容则是相反的过程。