4.3.3 Spark Streaming工作原理
1)Spark Sreaming计算流程
Spark Streaming可将流式计算分解成一系列短小的批处理作业。这里的批处理引擎是Spark Core,也就是把Spark Streaming的输入数据按照Batch Size(如1秒)分成一段一段的数据DStream(Discretized Stream),每一段数据都转换成Spark中的RDD,然后将Spark Streaming中对DStream的Transformation操作变为针对Spark中对RDD的Transformation操作,将RDD操作生成的中间结果保存在内存中。整个流式计算根据业务的需求可以对中间的结果进行叠加或者存储到外部设备。图4⁃13为Spark Streaming的整个计算流程。

图4⁃13 Spark Streaming计算流程
2)Spark Sreaming容错机制
对于流式计算来说,容错性至关重要。首先要明确一下Spark中RDD的容错机制。Spark中,每个RDD都是一个不可变的分布式可重算的数据集,其记录着确定性的操作继承关系(lineage),所以只要输入数据是可容错的,那么任意一个RDD的分区(Partition)出错或不可用,都是可以利用原始输入数据通过转换操作而重新算出。
对于Spark Streaming来说,其RDD的传承关系如图4⁃14所示,图中的每一个椭圆形都表示一个RDD,椭圆形中的每个圆形代表一个RDD中的一个Partition,图中的每一行的多个RDD表示一个DStream(图中有两个DStream),而每一行最后一个RDD则表示每一个Batch Size所产生的中间结果RDD。可以看到图中的每一个RDD都是通过血统(lineage)相连接的,由于Spark Streaming输入数据可以来自磁盘,例如HDFS(多份拷贝)或是来自网络的数据流(Spark Streaming会将网络输入数据的每一个数据流拷贝两份到其他的机器)都能保证容错性,所以RDD中任意的Partition出错,都可以并行地在其他机器上将缺失的Partition计算出来。这个容错恢复方式比连续计算模型(如Storm)的效率更高。

图4⁃14 Spark Streaming中RDD的继承关系
Spark Streaming将流式计算分解成多个Spark Job,对于每一段数据的处理都会经过Spark DAG图分解以及Spark的任务集的调度过程。对于目前版本的Spark Streaming而言,其最小的Batch Size的选取在0.5~2 s(Storm目前最小的延迟是100 ms左右),所以Spark Streaming能够满足除对实时性要求非常高(如高频实时交易)之外的所有流式准实时计算场景。
3)Spark Sreaming工作原理
Spark Streaming的基本工作原理是将输入数据流以时间片(秒级)为单位进行拆分,然后以类似批处理的方式处理每个时间片数据,其基本工作原理如图4⁃15所示。(https://www.daowen.com)

图4⁃15 Spark Streaming基本工作原理
首先,Spark Streaming把实时输入数据流以时间片
t (如1 s)为单位切分成块。Spark Streaming会把每块数据作为一个RDD,并使用RDD操作处理每一小块数据。每个块都会生成一个Spark Job处理,最终结果也返回多块。
使用Spark Streaming编写的程序与编写Spark程序非常相似,在Spark程序中,主要通过操作RDD提供的接口,如map、reduce、filter等,实现数据的批处理。而在Spark Streaming中,则通过操作DStream(表示数据流的RDD序列)提供的接口实现数据处理,这些接口和RDD提供的接口类似。图4⁃16和图4⁃17展示了由Spark Streaming程序到Spark jobs的转换图。

图4⁃16 Spark Streaming程序转换为DStream Graph
图4⁃17中,Spark Streaming把程序中对DStream的操作转换为DStream Graph。对于每个时间片,DStream Graph都会产生一个RDD Graph,针对每个输出操作(如print、foreach等),Spark Streaming都会创建一个Spark action,对于每个Spark action,Spark Streaming都会产生一个相应的Spark job,并交给JobManager。JobManager中维护着一个jobs队列, Spark job存储在这个队列中,JobManager把Spark job提交给Spark Scheduler,Spark Scheduler负责调度Task到相应的Spark Executor上执行,最后形成 Spark的job。

图4⁃17 DStream Graph转换为Spark jobs
正如Spark Streaming最初的目标一样,它通过丰富的API和基于内存的高速计算引擎让用户可以结合流式处理、批处理和交互查询等应用。因此Spark Streaming适合一些需要历史数据和实时数据结合分析的应用场合。当然,对于实时性要求不是特别高的应用也能完全胜任。另外通过RDD的数据重用机制可以得到更高效的容错处理。