4.3.4 Spark Streaming流计算编程模型
1)Spark Sreaming编程思想
DStream(Discretized Stream)作为Spark Streaming的基础抽象,它代表持续性的数据流。这些数据流既可以通过外部输入源来获取,也可以通过现有的DStream进行Transformation操作来获得。在内部实现上,DStream由一组时间序列上连续的RDD来表示。每个RDD都包含了自己特定时间间隔内的数据,如图4⁃18所示。

图4⁃18 DStream在时间轴下生成离散的RDD序列
对数据的操作也是以RDD为单位来进行的。在流数据分成一批一批后,生成一个先进先出的队列,然后Spark Engine从该队列中依次取出一个个批数据,把批数据封装成一个RDD,然后进行处理,如图4⁃19所示。

图4⁃19 DStream处理流程
作为构建于Spark之上的应用框架,Spark Streaming承袭了Spark的编程风格。
2)DStream输入源
在Spark Streaming中所有的操作都是基于流的,而输入源是这一系列操作的起点。输入DStreams和DStreams接收的流都代表输入数据流的来源,在Spark Streaming提供两种内置数据流来源。
(1)基础来源
在StreamingContext API中,它直接可用的来源。例如:文件系统、Socket(套接字)连接和Akka actors。
Spark Streaming提供了streamingContext.fileStream(dataDirectory)方法,可以从任何文件系统(如:HDFS、S3、NFS等)的文件中读取数据,然后创建一个DStream。Spark Streaming监控dataDirectory目录和在该目录下任何文件被创建处理(不支持在嵌套目录下写文件)。需要注意的是,读取的必须是具有相同数据格式的文件,创建的文件必须在dataDirectory目录下,并通过自动移动或重命名成数据目录。文件一旦移动就不能被改变,如果文件被不断追加,新的数据将不会被阅读。对于简单的文本文件,可以使用一个简单的方法streamingContext.textFileStream(dataDirectory)来读取数据。
Spark Streaming也可以基于自定义Actors的流创建DStream,通过Akka actors接受数据流,使用方法streamingContext.actorStream(actorProps, actor⁃name)。Spark Streaming使用streamingContext.queueStream(queueOfRDDs)方法可以创建基于RDD队列的DStream,每个RDD队列将被视为DStream中一块数据流进行加工处理。
(2)高级来源
可以通过额外的实用工具类来创建,如Kafka、Flume、Kinesis、Twitter等。
这一类的来源需要外部non⁃Spark库的接口,其中一些来源有复杂的依赖(如Kafka、Flume)。因此通过这些来源创建DStreams需要明确其依赖。例如,如果想创建一个使用Twitter tweets的数据的DStream流,必须按以下步骤处理。
•在SBT或Maven工程里添加spark⁃streaming⁃twitter_2.10依赖。
•开发:导入TwitterUtils包,通过TwitterUtils.createStream方法创建一个DStream。
•部署:添加所有依赖的jar包(包括依赖的spark⁃streaming⁃twitter_2.10及其依赖),然后部署应用程序。
需要注意的是,这些高级的来源一般在Spark Shell中不可用,因此基于这些高级来源的应用不能在Spark Shell中进行测试。如果必须在Spark Shell中使用它们,需要下载相应的Maven工程的Jar依赖并添加到类路径中。
需要重申的一点是在开始编写Spark Streaming程序之前,一定要将高级来源依赖的Jar添加到SBT或Maven项目相应的artifact中。常见的输入源和其对应的Jar包见表4⁃3。
表4⁃3 常见输入源的Jar包

另外,输入DStream也可以创建自定义的数据源,需要做的就是实现一个用户定义的接收器。
3)DStream操作(https://www.daowen.com)
与RDD类似,DStream也提供了自己的一系列操作方法,这些操作可以分成三类:普通转换操作、窗口转换操作和输出操作。
(1)普通转换操作
DStream的普通转换操作见表4⁃4。
表4⁃4 DStream普通转换操作


(2)窗口转换操作
DStream的窗口转换操作见表4⁃5。
表4⁃5 DStream窗口转换操作


在Spark Streaming中,数据处理是按批进行的,而数据采集是逐条进行的,因此在Spark Streaming中会先设置好批处理间隔(Batch Duration),当超过批处理间隔的时候就会把采集的数据汇总起来成为一批数据交给系统去处理。
对于窗口操作而言,在其窗口内部会有N个批处理数据,批处理数据的大小由窗口间隔(Window Duration)决定,而窗口间隔指的就是窗口的持续时间,在窗口操作中,只有窗口的长度满足了才会触发批数据的处理。除了窗口的长度,窗口操作还有另一个重要的参数就是滑动间隔(Slide Duration),它指的是经过多长时间窗口滑动一次形成新的窗口,滑动窗口默认情况下和批次间隔的相同,而窗口间隔一般设置得要比它们两个大。在这里必须注意的一点是滑动间隔和窗口间隔的大小一定得设置为批处理间隔的整数倍。
如图4⁃20所示,批处理间隔是1个时间单位,窗口间隔是3个时间单位,滑动间隔是2个时间单位。对于初始的窗口time 1-time 3,只有窗口间隔满足了才触发数据的处理。这里需要注意的一点是,初始的窗口有可能流入的数据没有撑满,但是随着时间的推进,窗口最终会被撑满。当每隔2个时间单位,窗口滑动一次后,会有新的数据流入窗口,这时窗口会移去最早的两个时间单位的数据,而与最新的两个时间单位的数据进行汇总形成新的窗口(time 3-time 5)。

图4⁃20 批处理间隔示意图
对于窗口操作,批处理间隔、窗口间隔和滑动间隔是非常重要的三个时间概念,是理解窗口操作的关键所在。
(3)输出操作
Spark Streaming允许DStream的数据被输出到外部系统,如数据库或文件系统。由于输出操作实际上使转换操作后的数据可以通过外部系统被使用,同时输出操作触发所有DStream的转换操作的实际执行(类似于RDD操作)。
DStream主要的输出操作见表4⁃6。
表4⁃6 DStream输出操作
