Storm-runtime

master-slave结构:
- Nimbus是主节点,负责分发用户代码,指派Supervisor上的worker进程,运行topology的(Spout/Bolt)Task
- Supervisor是从节点,守护进程. 负责启动和终止worker进程. 通过Storm的配置文件中的 supervisor.slots.ports配置项,可以指定在一个Supervisor上最大允许多少个Slot,每个Slot通过端口号来唯一标识,一个端口号 对应一个Worker进程(如果该Worker进程被启动)。

运行流程
1)户端提交拓扑到nimbus。
2) Nimbus针对该拓扑建立本地的目录根据topology的配置计算task,分配task,在zookeeper上建立assignments节点存储 task和supervisor机器节点中woker的对应关系;
在zookeeper上创建taskbeats节点来监控task的心跳;启动topology。
3) Supervisor去zookeeper上获取分配的tasks,启动多个woker进行,每个woker生成task,一个task一个线程;根据topology 信息初始化建立task之间的连接;Task和Task之间是通过zeroMQ管理的;后整个拓扑运行起来。
内核-队列
在Storm中大量使用Disrupt Queue来解耦Storm内部的消息处理过程。分析这些Queue的分布,是分析Storm运行时态的基础。

在Storm的worker中,最小的执行单元是executor,一个executor只会有一个component。目前一个component只会有一个task。而一个executor会有两个Disruptor Queue, 一个用于接受数据的receive Disrupor Queue和一个用于发送send Disrupor Queue。然后在worker中有一个全局的send Disrupor Queue。分析完队列的分布,在分析topology的运行时状况。
当一个Topology提交到Storm集群后,task被分配到各个wrker中开始执行后,task是怎么执行的。Storm的Task分为两类,一类是消息的源头Spout,它负责在源头产生消息;然后就是Bolt,它是执行单元。但是这两者各自的工作逻辑如下。
我们来分析Spout在运行时的工作状况。Spout的nextTuple用于发送数据,接口注释上说明它不能阻塞,因为它与active, deactive,ack, fail在一个处理线程里面被处理。但是nextTuple被阻塞会有什么副作用列?当nextTuple被阻塞,应用代码中另起线程调用SpoutCollector会有什么后果?下图是SpoutExecutor的执行逻辑。

在SpoutExecutor处理循环里面,第一步做的事情是从receive Disrupor Queue里面消费里面的消息。这里面的就是SpoutExecutor所收到的消息。处理逻辑如下图所示

然后SpoutExecutor会将overflow中的数据再次发送。overflow是用于接受SpoutCollector.emit()无发及时发送的数据的,具体SpoutCollector发送数据的逻辑见下文分析。但是这里在发送overflow的数据时,与SpoutCollector.emit()有个区别,就是数据还是无法被正常发送时,数据会丢弃,也就是不会被再次写入overflow中。当overflow中没有数据,以及pending中的数据量小于TOPOLOGY_MAX_SPOUT_PENDING时,判断topology的状态。当topology不是deactive状态时,如果topology有触发active命令,会调用spout的active接口,然后调用spout的nextTuple接口。否则调用deactive接口。所以这里当Spout的nextTuple被阻塞时,spout没办法处理acker回报的消息,回报的消息会阻塞在executor的receive DisruptorQueue中,当receive DisruptorQueue塞满后,是会阻塞在对应的网络处理模块中,storm中经典的是zeroMQ,老版本的zeroMQ是没有设置水位,这样会大量堆积到内存中,因为zeroMQ是C++的,占用的是堆外内存,JVM无法管理,最坏就是把机器的内存耗光。如果这个是另起线程调用SpoutCollector的emit,当emit数据速度过快,会导致overflow中堆积数据,导致worker内存消耗。
上面讲到SpoutCollector在emit数据会写入overflow,那什么情况下会写入overflow。SpoutCollector.emit是否真的就把数据发送到了网络。下面是SpoutCollector的处理逻辑。

Spout当调用SpoutCollect.emit()发送时,首先的逻辑是根据根据消息的STREAM ID,和对应的values,得到目的端task id。当topology开启acker机制时,会生成一个随机的rootId,否则使用一个默认值作为rootId。然后为消息生成一个随机的messageId,由rootId和messageId组成对应的tuple Id。这里tuple Id是有rootId为key,messageId为value的一个HashMap。这里采用HashMap的作用会在下面Spout的ack机制是做说明。在将生成的tuple Id和用会的values组成一个tuple。并判断overflower是否有数据,让overflow中有数据时,数据会直接放入overflow中,而不会放入executor 的send Disraptor Queue中。当overflow为空时,会将tuple放入send Disraptor Queue中。当捕获Disraptor Queue的InsufficientCapacityException时,数据就放入了overflow中。然后处理ack。当开启了ack机制,会在pending中加入对应的tuple数据,然后项acker发送ack init消息。
根据上面的发送过程,数据emit只是被写入了executor的send Disraptor Queue。而数据在Spout端的丢失多是数据在overflow中被SpoutExecutor在次处理时,send Disraptor Queue满导致。
关于pending,首先它是一个RotatingMap。它通过定时旋转,以达到定时器的目的。在SpoutExecutor的RotatingMap中有两个桶,然后executor有个定时线程,会按照用户设定的message timeout second,定期向executor的receive DisruptorQueue写入SYSTEM_TICK_STREAM_ID消息,然后SpoutExecutor处理SYSTEM_TICK_STREAM_ID消息时就旋转RotatingMap,当被清除的桶中有数据时,被清除的桶中的数据会调用fail接口,通知业务逻辑fail。这就是当超时间设置比业务逻辑短时,导致数据重复的原因。还有一种是acker消息丢失,导致数据重复。由于RotatingMap的底层是HashMap,中间没有锁,overflow是个LinkedList,所以SpoutCollector和SpoutExecutor不能在两个线程中并发 处理。
BoltExecutor的处理逻辑非常简单,就是消费receive Disrupor Queue中数据,然后调用bolt的excue。但是这里也是单线程处理的,阻塞或者处理速度不匹配,就会导致数据在Disrupor Queue或者网络模块中堆积,其中使用zeroMQ的副作用最大。
BoltCollector的emit与SpoutCollector的emit处理相比,首先是少了overflow承接无法发送的数据,会直接丢弃。其次是没有pending和acker消息的发送。
接下来分析一下Storm的Acker机制。Storm的Ack机制在Storm刚刚开源时被大书特书,处理原理也确实非常的精彩。由于这里ID都是随机数,所以这里不会在讨论随机数的唯一问题。

如上图所示,spout发送t1给bolt a, bolt a 在t1的基础上生成t2,t3,t4给bolt b,bolt b ack所收到的数据。下面来追踪整个id变化的过程。
Spout 发送t1, message id为<r_1, m_1>,发送ack
Acker收到ack init, map中缓存<r_1,m_1>
Bolt a 收到t1, messageId为<r_1,m_1>,生成t2,t3, t4
Bolt a 发送t2, anchor t1, t2,messageId为<r_1,m_2>, 更新t1的ackVal为m_2
Bolt a 发送t3, anchor t1, t3,messageId为<r_1,m_3>, 更新t1的ackVal为m_2^m_3
Bolt a 发送t4, anchor t1, t4,messageId为<r_1,m_4>, 更新t1的ackVal为m_2^m_3^m_4
Bolt a ack t1, 向acker发送ack消息<r_1, m_1^ m_2^m_3^m_4>
Acker 收到bolt a的ack消息,更新缓存为<r_1, m_1^ m_1^ m_2^m_3^m_4>即<r_1, m_2^m_3^m_4>
Bolt b ack t2, 向acker发送ack消息<r_1, m_2>
Acker 收到bolt b的ack消息,更新缓存为<r_1, m_2^ m_2^m_3^m_4>即<r_1, m_3^m_4>
Bolt b ack t3, 向acker发送ack消息<r_1, m_3>
Acker 收到bolt b的ack消息,更新缓存为<r_1, m_3^m_3^m_4>即<r_1, m_4>
Bolt b ack t4, 向acker发送ack消息<r_1, m_4>
Acker 收到bolt b的ack消息,更新缓存为<r_1, m_4^m_4>即<r_1, 0>
messageId确认完毕,向Spout发送ack消息。当消息没有被ack,会一直在spout的pending队列中,知道被ack或者超时。
它基本上使用两个long值就跟踪了一个消息在整个流中的处理过程。
【参考文献】