本文仅提供一个入门概览,部分内容来源于网络,部分来源于自己理解,参考内容链接会在文末给出,部分内容未找到原作,如有侵权,请联系删除。

1、概述

许多分布式计算系统都可以实时或者接近实时地处理大数据流。Storm是一个免费并开源的分布式实时计算系统。利用Storm可以很容易做到可靠地处理无限的数据流,像Hadoop批量处理大数据一样,Storm可以实时处理数据。Hadoop 在本质上是一个批处理系统。数据被引入Hadoop 文件系统(HDFS) 并分发到各个节点进行处理。当处理完成时,结果数据返回到 HDFS 供始发者使用。Storm 支持创建拓扑结构来转换没有终点的数据流。不同于 Hadoop 作业,这些转换从不停止,它们会持续处理到达的数据。

2、Storm相关基础概念

在Storm的集群里面有两种节点:控制节点(master node)和工作节点(worker node)。控制节点上面运行一个后台程序:Nimbus,它的作用类似Hadoop里面的JobTracker。Nimbus负责在集群里面分布代码,分配工作给机器,并且监控状态。每一个工作节点上面运行一个叫做Supervisor的节点(类似 TaskTracker)。Supervisor会监听分配给它那台机器的工作,根据需要启动/关闭工作进程。每一个工作进程执行一个Topology(类似 Job)的一个子集;一个运行的Topology由运行在很多机器上的很多工作进程Worker(类似Child)组成。

img

 

3、storm与hadoop的比较

img

Storm 与Hadoop最大的不同之处在于它的处理方式。Hadoop 在本质上是一个批处理系统。数据被引入Hadoop 文件系统(HDFS) 并分发到各个节点进行处理。当处理完成时,结果数据返回到 HDFS 供始发者使用。Storm 支持创建拓扑结构来转换没有终点的数据流。不同于 Hadoop 作业,这些转换从不停止,它们会持续处理到达的数据。

你也可以将它们形象地理解为桶装水与自来水的区别。Hadoop需要一份一份地进行打包、处理、搬运,通常来说,从“水源水”到“饮用水”需要等待比较长的一段时间。Storm则建立了一套管网系统,管网系统中可以任意地添加“水”处理环节,或者安装“水龙头”进行数据输入和输出,从而整个处理过程变得更加流畅、快捷。Hadoop和Storm框架用于分析大数据,两者互补,在某些方面有所不同。Apache Storm执行除持久性之外的所有操作,而Hadoop在所有方面都很好,但滞后于实时计算。可以总结以下几点:

1、Storm用于实时计算,Hadoop用于离线计算。

2、Storm处理的数据保存在内存中,源源不断;

3、Hadoop处理的数据保存在文件系统中,一批一批。

4、Storm的数据通过网络传输进来;Hadoop的数据保存在磁盘中。

5、Storm与Hadoop的编程模型相似

下面的表格是W3C school上比较了Storm和Hadoop属性的结果(原文在这里:https://www.w3cschool.cn/apache_storm/apache_storm_introduction.html)。

StormHadoop
实时流处理批量处理
无状态有状态
主/从架构与基于ZooKeeper的协调。主节点称为nimbus,从属节点是主管。具有/不具有基于ZooKeeper的协调的主 - 从结构。主节点是作业跟踪器,从节点是任务跟踪器。
Storm流过程在集群上每秒可以访问数万条消息。Hadoop分布式文件系统(HDFS)使用MapReduce框架来处理大量的数据,需要几分钟或几小时。
Storm拓扑运行直到用户关闭或意外的不可恢复故障。MapReduce作业按顺序执行并最终完成。
两者都是分布式和容错的 
如果nimbus / supervisor死机,重新启动使它从它停止的地方继续,因此没有什么受到影响。如果JobTracker死机,所有正在运行的作业都会丢失。

 

4、Storm集群

1、Storm集群架构

Storm集群遵循主/从结构。Storm的主节点是半容错的。Strom集群由一个主节点(nimbus)和一个或者多个工作节点(supervisor)组成。 除此之外Storm集群还需要一个ZooKeeper的来进行集群协调。架构如下图所示:

2、相关术语

Nimbus

Storm集群的Master节点,负责分发用户代码,指派给具体的Supervisor节点上的Worker节点,去运行Topology对应的组件(Spout/Bolt)的Task。

Supervisor

Storm集群的从节点,负责管理运行在Supervisor节点上的每一个Worker进程的启动和终止。通过Storm的配置文件中的supervisor.slots.ports配置项,可以指定在一个Supervisor上最大允许多少个Slot,每个Slot通过端口号来唯一标识,一个端口号对应一个Worker进程(如果该Worker进程被启动)。

Worker

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

Task

worker中每一个spout/bolt的线程称为一个task. 在storm0.8之后,task不再与物理线程对应,不同spout/bolt的task可能会共享一个物理线程,该线程称为executor。

ZooKeeper

用来协调Nimbus和Supervisor,如果Supervisor因故障出现问题而无法运行Topology,Nimbus会第一时间感知到,并重新分配Topology到其它可用的Supervisor上运行、

5、Storm编程模型

Strom在运行中可分为spout与bolt两个组件,其中,数据源从spout开始,数据以tuple的方式发送到bolt,多个bolt可以串连起来,一个bolt也可以接入多个spot/bolt。运行时原理如下图:

 

img

从应用上来说,在Storm中,需要先设计一个实时计算结构,我们称之为拓扑(topology)。之后,这个拓扑结构会被提交给集群,其中主节点(master node)负责给工作节点(worker node)分配代码,工作节点负责执行代码。在一个拓扑结构中,包含spout和bolt两种角色。数据在spouts之间传递,这些spouts将数据流以tuple元组的形式发送;而bolt则负责转换数据流。

img

1、Topologies

为了在storm上面做实时计算,需要建立一些图状结构,我们称之为(topologies)。一个topology是spouts和bolts组成的图,其中spout负责发送消息,负责将数据流以tuple元组的形式发送出去;而bolt则负责转换这些数据流,在bolt中可以完成计算、过滤等操作,bolt自身也可以随机将数据发送给其他bolt。

Spouts和bolts通过stream groupings(定义一个流在Bolt任务间该如何被切分)连接起来,我们也可以将它理解为一个由无限制的处理节点组成的图状结构。Topology里面的每个处理节点都包含处理逻辑,而节点之间的连接则表示数据流动的方向。

2、Stream groupings

定义一个topology的其中一步是定义每个bolt接收什么样的流作为输入。stream grouping就是用来定义一个stream应该如果分配数据给bolts上面的多个tasks。

img

 

3、Streams

消息流stream是storm里的关键抽象。一个消息流是一个没有边界的tuple序列( tuple是一个类似于列表的东西,存储的每个元素叫做field),而这些tuple序列会以一种分布式的方式并行地创建和处理。通过对stream中tuple序列中每个字段命名来定义stream。

4、Spouts

消息源spout是Storm里面一个topology里面的消息生产者。一般来说消息源会从一个外部源读取数据并且向topology里面发出消息:tuple。Spout可以是可靠的也可以是不可靠的。如果这个tuple没有被storm成功处理,可靠的消息源spouts可以重新发射一个tuple,但是不可靠的消息源spouts一旦发出一个tuple就不能重发了。

5、Bolts

所有的消息处理逻辑被封装在bolts里面。Bolts可以做很多事情:过滤,聚合,查询数据库等等, Bolts也可以简单的做消息流的传递。复杂的消息流处理往往需要很多步骤,从而也就需要经过很多bolts。比如算出一堆图片里面被转发最多的图片就至少需要两步:第一步算出每个图片的转发数量。第二步找出转发最多的前10个图片。(如果要把这个过程做得更具有扩展性那么可能需要更多的步骤)。

img

6、以上概念的一句话概括

Topology:Storm中运行的一个实时应用程序的名称。将 Spout、 Bolt整合起来的拓扑图。定义了 Spout和 Bolt的结合关系、并发数量、配置等等。

Spout:在一个topology中获取源数据流的组件。通常情况下spout会从外部数据源中读取数据,然后转换为topology内部的源数据。

Bolt:接受数据然后执行处理的组件,用户可以在其中执行自己想要的操作。

Tuple:一次消息传递的基本单元,理解为一组消息就是一个Tuple。

Stream:Tuple的集合。表示数据的流向。

(备注:详解见中文手册,文末有链接)

6、并发机制

1、并发级别

Nodes

服务器:配置在Storm集群中的一个服务器,会执行Topology的一部分运算,一个Storm集群中包含一个或者多个Node。

Workers

JVM虚拟机、进程:指一个Node上相互独立运作的JVM进程,每个Node可以配置运行一个或多个worker。一个Topology会分配到一个或者多个worker上运行。

Executor

线程:指一个worker的jvm中运行的java线程。多个task可以指派给同一个executer来执行。除非是明确指定,Storm默认会给每个executor分配一个task。

Task

bolt/spout实例:task是spout和bolt的实例,他们的nextTuple()和execute()方法会被executors线程调用执行。

大多数情况下,除非明确指定,Storm的默认并发设置值是1。即,一台服务器(node),为topology分配一个worker,每个executer执行一个task。如图:Storm默认并发机制,此时唯一的并发机制出现在线程级即Executor。

img

2、增加各级别并发

增加Node

这个其实就是增加集群的服务器数量。

增加worker

可以通过API和修改配置两种方式修改分配给topology的woker数量。

API增加woker:

Config config = new Config();
config.setNumWorkers(2);

单机模式下,增加worker的数量不会有任何提升速度的效果。

增加Executor

API增加Executor:

builder.setSpout(spout_id,spout,2);
builder.setBolt(bolt_id,bolt,executor_num);

这种办法为Spout或Bolt增加线程数量,默认每个线程都运行该Spout或Bolt的一个task。

增加Task

API增加Task:

builder.setSpout(...).setNumTasks(2);
builder.setBolt(...).setNumTasks(task_num);

如果手动设置过task的数量,task的总数量就是指定的数量个,而不管线程有几个,这些task会随机分配在这些个线程内部执行。

3、数据流分组

数据流分组方式定义了数据如何进行分发。

Storm内置了七种数据流分组方式:

Shuffle Grouping

随机分组。

随机分发数据流中的tuple给bolt中的各个task,每个task接收到的tuple数量相同。

Fields Grouping

按字段分组。

根据指定字段的值进行分组。指定字段具有相同值的tuple会路由到同一个bolt中的task中。

All Grouping

全复制分组。

所有的tuple复制后分发给后续bolt的所有的task。

Globle Grouping

全局分组。

这种分组方式将所有的tuple路由到唯一一个task上,Storm按照最小task id来选取接受数据的task。这种分组方式下配置bolt和task的并发度没有意义。

这种方式会导致所有tuple都发送到一个JVM实例上,可能会引起Strom集群中某个JVM或者服务器出现性能瓶颈或崩溃。

None Grouping

不分组。

在功能上和随机分组相同,为将来预留。

Direct Grouping

指向型分组。

数据源会通过emitDirect()方法来判断一个tuple应该由哪个Strom组件来接受。只能在声明了是指向型数据流上使用。

Local or shuffle Grouping

本地或随机分组。

和随机分组类似,但是,会将tuple分发给同一个worker内的bolt task,其他情况下采用随机分组方式。这种方式可以减少网络传输,从而提高topology的性能。

自定义

另外可以自定义数据流分组方式

写类实现CustomStreamGrouping接口

代码:

/**
* 自定义数据流分组方式
* @author park
*
*/
public class MyStreamGrouping implements CustomStreamGrouping {
​
  /**
  * 运行时调用,用来初始化分组信息
  * context:topology上下文对象
  * stream:待分组数据流属性
  * targetTasks:所有待选task的标识符列表
  * 
 */
  @Override
  public void prepare(WorkerTopologyContext context, GlobalStreamId stream, List<Integer> targetTasks) {
​
  }
​
  /**
  * 核心方法,进行task选择
  * taskId:发送tuple的组件id
  * values:tuple的值
  * 返回值要发往哪个task
  */
  @Override
  public List<Integer> chooseTasks(int taskId, List<Object> values) {
    return null;
  }
}

7、strom的优势

简单编程 在大数据处理方面相信大家对hadoop已经耳熟能详,基于GoogleMap/Reduce来实现的Hadoop为开发者提供了map、reduce原语,使并行批处理程序变得非常地简单和优美。同样,Storm也为大数据的实时计算提供了一些简单优美的原语,这大大降低了开发并行实时处理的任务的复杂性,帮助你快速、高效的开发应用。 Spark提供的数据集操作类型有很多种,不像Hadoop只提供了Map和Reduce两种操作。比如map,filter, flatMap,sample, groupByKey, reduceByKey, union, join, cogroup,mapValues, sort,partionBy等多种操作类型,他们把这些操作称为Transformations。同时还提供Count, collect, reduce, lookup, save等多种actions。这些多种多样的数据集操作类型,给上层应用者提供了方便。各个处理节点之间的通信模型不再像Hadoop那样就是唯一的Data Shuffle一种模式。用户可以命名,物化,控制中间结果的分区等。可以说编程模型比Hadoop更灵活.。

多语言支持 除了用java实现spout和bolt,你还可以使用任何你熟悉的编程语言来完成这项工作,这一切得益于Storm所谓的多语言协议。多语言协议是Storm内部的一种特殊协议,允许spout或者bolt使用标准输入和标准输出来进行消息传递,传递的消息为单行文本或者是json编码的多行。 Storm支持多语言编程主要是通过ShellBolt,ShellSpout和ShellProcess这些类来实现的,这些类都实现了IBolt 和 ISpout接口,以及让shell通过java的ProcessBuilder类来执行脚本或者程序的协议。 可以看到,采用这种方式,每个tuple在处理的时候都需要进行json的编解码,因此在吞吐量上会有较大影响。

支持水平扩展 在Storm集群中真正运行topology的主要有三个实体:工作进程、线程和任务。Storm集群中的每台机器上都可以运行多个工作进程,每个工作进程又可创建多个线程,每个线程可以执行多个任务,任务是真正进行数据处理的实体,我们开发的spout、bolt就是作为一个或者多个任务的方式执行的。 因此,计算任务在多个线程、进程和服务器之间并行进行,支持灵活的水平扩展。

容错性强 如果在消息处理过程中出了一些异常,Storm会重新安排这个出问题的处理单元。Storm保证一个处理单元永远运行(除非你显式杀掉这个处理单元)。

可靠的消息保证 Storm可以保证spout发出的每条消息都能被“完全处理”,这也是直接区别于其他实时系统的地方,如S4。

快速的消息处理 用ZeroMQ作为底层消息队列, 保证消息能快速被处理

本地模式,支持快速编程测试 Storm有一种“本地模式”,也就是在进程中模拟一个Storm集群的所有功能,以本地模式运行topology跟在集群上运行topology类似,这对于我们开发和测试来说非常有用。

8、总结

最后再来梳理一下Storm中涉及的主要概念:

1.拓扑(Topology):打包好的实时应用计算任务,同Hadoop的MapReduce任务相似。

2.元组(Tuple):是Storm提供的一个轻量级的数据格式,可以用来包装你需要实际处理的数据。

3.流(Streams):数据流(Stream)是Storm中对数据进行的抽象,它是时间上无界的tuple元组序列(无限的元组序列)。

4.Spout(喷嘴):Storm中流的来源。Spout从外部数据源,如消息队列中读取元组数据并吐到拓扑里。

5.Bolts:在拓扑中所有的计算逻辑都是在Bolt中实现的。

6.任务(Tasks):每个Spout和Bolt会以多个任务(Task)的形式在集群上运行。

7.组件(Component):是对Bolt和Spout的统称。

8.流分组(Stream groupings):流分组定义了一个流在一个消费它的Bolt内的多个任务(task)之间如何分组。

9.可靠性(Reliability):Storm保证了拓扑中Spout产生的每个元组都会被处理。

10.Workers(工作进程):拓扑以一个或多个Worker进程的方式运行。每个Worker进程是一个物理的Java虚拟机,执行拓扑的一部分任务。

11.Executor(线程):是1个被worker进程启动的单独线程。每个executor只会运行1个topology的1个component。

12.Nimbus:Storm集群的Master节点,负责分发用户代码,指派给具体的Supervisor节点上的Worker节点,去运行Topology对应的组件(Spout/Bolt)的Task。

13.Supervisor:Storm集群的从节点,负责管理运行在Supervisor节点上的每一个Worker进程的启动和终止

 

 

以下为参考文档:

https://zhuanlan.zhihu.com/p/25414207:实时大数据顶级分析工具Storm图文解析

https://www.cnblogs.com/gzxbkk/p/9456029.html:storm架构及原理

https://blog.csdn.net/u011082453/article/details/82417259:storm简介、原理、概念

https://www.cntofu.com/book/108/doc/zh/Concepts.md:storm官方中文手册

https://my.oschina.net/u/3754001/blog/1805746:Storm介绍及原理

https://zhuanlan.zhihu.com/p/28658526:一文让你重拾所有的Storm架构与运行原理

https://zhuanlan.zhihu.com/p/64112857:大数据Storm相比于Spark、Hadoop有哪些优势

https://www.cnblogs.com/panfeng412/tag/Storm/:大圆那些事

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐