大数据学习笔记(三):Storm
本文仅提供一个入门概览,部分内容来源于网络,部分来源于自己理解,参考内容链接会在文末给出,部分内容未找到原作,如有侵权,请联系删除。
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)组成。

3、storm与hadoop的比较

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)。
| Storm | Hadoop |
|---|---|
| 实时流处理 | 批量处理 |
| 无状态 | 有状态 |
| 主/从架构与基于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。运行时原理如下图:

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

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。

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个图片。(如果要把这个过程做得更具有扩展性那么可能需要更多的步骤)。

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。

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有哪些优势
更多推荐
所有评论(0)