4.2.2 Storm流计算架构

更新于 2026年10月10日 版权声明
4.2.2 Storm流计算架构

1)Storm集群架构

Storm集群采用主从架构方式,主节点是Nimbus,从节点是Supervisor,有关调度相关的信息存储到Zookeeper集群中,Storm集群架构如图4⁃3所示。

①Nimbus:负责资源分配和任务调度,负责分发用户代码,指派给具体的Supervisor节点上的Worker节点,去运行Topology对应的组件(Spout/Bolt)的Task。

②Supervisor:负责接收Nimbus分配的任务,启动和停止属于自己管理的Worker进程。

③Worker:负责运行具体处理组件逻辑的进程。Worker运行的任务类型只有两种,一种是Spout任务,一种是Bolt任务。

④Zookeeper:用来协调Nimbus和Supervisor,存放如心跳信息、集群状态与配置信息等公有数据,Nimbus将分配给Supervisor的任务写入Zookeeper。

默认情况下,Storm集群有两种模式:本地模式和集群模式。

图示

图4⁃3 Storm集群架构

本地模式:此模式下,Storm Topology运行在本地机器的单个JVM上。此模式主要用于开发、测试和调试,因为它是看到所有的Topology组件一起工作的最简单的方法。在这种模式下,用户可以调整参数,并能够看到Topology在不同的Storm配置环境中的运行状况。

集群模式:在这种模式下,用户提交Topology到正常工作的Storm集群上,该集群由许多进程组成,通常运行在不同的机器上。

2)术语

(1)拓扑(Topology)

一个Storm拓扑打包了一个实时处理程序的逻辑,因为各个组件间的消息流动形成逻辑上的一个拓扑结构。一个Storm拓扑跟一个MapReduce的任务(Job)是类似的,主要区别是MapReduce任务最终会结束,而拓扑任务会一直运行(当然直到杀死它)。

(2)元组(Tuple)

元组是Storm提供的一个轻量级的数据格式,可以用来包装用户需要实际处理的数据。元组是一次消息传递的基本单元。一个元组是一个命名的值列表,其中的每个值都可以是任意类型的。元组是动态地进行类型转化的,即字段的类型不需要事先声明。在Storm中编程时,就是在操作和转换由元组组成的流。通常,元组包含整数、字节、字符串、浮点数、布尔值和字节数组等类型。开发人员要想在元组中使用自定义类型,就需要实现元组自己的序列化方式。

(3)流(Stream)

流是Storm中的核心抽象。一个流由无限的元组序列组成,这些元组会被分布式并行地创建和处理。通过流中元组包含的字段名称来定义这个流。每个流声明时都被赋予了一个ID。

(4)消息源(Spout)

Spout(喷嘴,这个名字很形象),在一个Topology中产生源数据流的组件。通常情况下Spout会从外部数据源中读取数据,然后转换为Topology内部的源数据,如消息队列中读取元组数据并吐到拓扑里。Spout可以是可靠的(reliable)或者不可靠(unreliable)的。可靠的Spout能够在一个元组被Storm处理失败时重新进行处理,而非可靠的Spout只是吐数据到拓扑里,不关心处理是成功了还是失败了。

Spout是一个主动的角色,其接口中有个nextTuple函数,Storm框架会不断调用它去做元组的轮询。如果没有新的元组过来,就会直接返回,否则把新元组吐到拓扑里。nextTuple必须是非阻塞的,因为Storm在同一个线程里执行Spout的函数。

Spout可以一次给多个流吐数据。此时就需要通过OutputFieldsDeclarer的declareStream函数来声明多个流,并在调用SpoutOutputCollector提供的emit方法时指定元组吐给哪个流。(https://www.daowen.com)

Spout中另外两个主要的函数是ack和fail。当Storm检测到一个从Spout吐出的元组在拓扑中成功处理完时调用ack,没有成功处理完时调用fail。只有可靠型的Spout会调用ack和fail函数。

(5)消息处理者(Bolt)

Bolt是在一个Topology中接受数据然后执行处理的组件。Bolt可以执行过滤、函数操作、合并、写数据库等任何操作。在拓扑中所有的计算逻辑都是在Bolt中实现的。Bolt是一个被动的角色,其接口中有个execute(Tuple input)函数,在接收到消息后会调用此函数,用户可以在其中执行自己想要的操作。一个Bolt可以处理任意数量的输入流,产生任意数量新的输出流。Bolt就是流水线上的一个处理单元,把数据的计算处理过程合理地拆分到多个Bolt,合理设置Bolt的task数量,能够提高Bolt的处理能力,提升流水线的并发度。

Bolt可以给多个流吐出元组数据。此时需要使用OutputFieldsDeclarer类的declareStream方法来声明多个流,使用OutputCollector对象来吐出新的元组,并在使用OutputColletor的emit方法时指定给哪个流吐数据。Bolt必须为处理的每个元组调用OutputCollector的ack方法,以便于Storm知道元组什么时候被各个Bolt处理完了(最终就可以确认Spout吐出的某个元组处理完了)。通常处理一个输入的元组时,会基于这个元组吐出零个或者多个元组,然后确认(ack)输入的元组处理完了,Storm提供了IBasicBolt接口来自动完成确认。必须注意OutputCollector不是线程安全的,所以所有的吐数据(emit)、确认(ack)、通知失败(fail)必须发生在同一个线程里。

当声明了一个Bolt的输入流,也就订阅了另外一个组件的某个特定的输出流。如果希望订阅另一个组件的所有流,需要单独挨个订阅。

(6)组件(Component)

在Storm中,组件是对Bolt和Spout的统称。

(7)流分组(Stream Grouping)

定义拓扑的时候,一部分工作是指定每个Bolt应该消费哪些流。流分组定义了一个流在一个消费它的Bolt内的多个任务(Task)之间如何分组。流分组跟计算机网络中的路由功能是类似的,决定了每个元组在拓扑中的处理路线。

在Storm中有七个内置的流分组策略,可以通过实现CustomStreamGrouping接口来自定义一个流分组策略。

(8)任务(Task)

Worker中每一个Spout/Bolt的线程称为一个Task。每个Spout和Bolt会以多个任务(Task)的形式在集群上运行。每个任务对应一个执行线程,流分组定义了如何从一组任务(同一个Bolt)发送元组到另外一组任务(另外一个Bolt)上。可以在调用TopologyBuilder的setSpout和setBolt函数时设置每个Spout和Bolt的并发数。

(9)轻量级消息内核(ZeroMQ,Zero Message Queue)

ZeroMQ(ZMQ)是一个基于消息队列的多线程网络库,其对套接字类型、连接处理、帧,甚至路由的底层细节进行抽象,提供跨越多种传输协议的套接字。也可以说,ZeroMQ是一个简单好用的传输层,是个类似于Socket的一系列接口。它使得Socket编程更加简单、简洁和性能更高。ZeroMQ跟Socket的区别是:普通的Socket是端到端(1:1)的关系,而ZeroMQ却可以是N:M的关系。人们对BSD套接字的了解较多的是点对点的连接,点对点连接需要显式地建立连接、销毁连接、选择协议(TCP/UDP)和处理错误等,而ZeroMQ屏蔽了这些细节,让网络编程更为简单。ZeroMQ是一个消息处理队列库,可在多个线程、内核和主机盒之间弹性伸缩。ZeroMQ的明确目标是“成为标准网络协议栈的一部分,之后进入Linux内核”。

3)Storm流计算模型

Storm实现了一个数据流的模型,在这个模型中数据持续不断地流经一个由很多转换实体构成的网络。一个数据流的抽象叫作流,流是无限的元组(Tuple)序列。元组就像一个可以表示标准数据类型(例如int,float和byte数组)和用户自定义类型的数据结构。每个流由一个唯一的ID来标识,这个ID可以用来构建拓扑中各个组件的数据源。

Storm流计算模型图如图4⁃4所示,其中的水龙头(Spout)代表了数据流的来源,一旦水龙头打开,数据就会源源不断地流经Bolt而被处理。图4⁃4中有三个流,每个数据流中流动的是元组(Tuple),它承载了具体的数据。元组通过流经不同的Bolt而被处理。

图示

图4⁃4 Storm流计算模型

在Hadoop中,数据的输入输出都需要放到自己的文件系统HDFS里。而Storm可以使用任意来源的数据输入和任意的数据输出,只要实现对应的代码来获取/写入这些数据就可以。输入/输出数据可以是基于类似Kafka或者ActiveMQ这样的消息队列,也可以是数据库、文件系统等。

↑上一章 ↓下一章
关注公众号获取验证码
复制内容需要验证码(7.99元/天)