4.3.2 Spark Streaming数据流
Spark Streaming是Spark核心API的一个扩展,是Spark用于实现实时流计算的核心组件。如图4⁃11所示,Spark Streaming支持从多种数据源获取数据,包括Kafka、Flume、HDFS/S3、Kinesis以及Twitter等,从数据源获取数据之后,可以使用诸如map、reduce、join和window等高级函数进行复杂算法的处理。最后还可以将处理结果存储到文件系统、数据库和现场仪表盘。

图4⁃11 Spark Streaming数据流图
1)Spark Streaming数据流处理过程
Spark Streaming在内部的处理机制是,接收实时流的数据,并根据一定的时间间隔拆分成一批批的数据,然后通过Spark Engine处理这些离散化的批数据,最终得到处理后的一批批结果数据,如图4⁃12所示。

图4⁃12 Spark Streaming数据流处理
离散后的一批数据在Spark内核中对应一个RDD实例。因此,对应流数据的DStream可以看成是一组RDD,即RDD的一个序列。也就是说,在流数据分成一批一批后,通过一个先进先出的队列,然后Spark Engine从该队列中依次取出一个个批数据DStream,并把批数据封装成一个RDD,最后进行处理。这是一个典型的生产者消费者模型,对应的就有生产者消费者模型的问题,即如何协调生产速率和消费速率。
输入的数据流经过Spark Streaming的Receive,数据切分为一个个离散的DStream(DStream是Spark Streaming中流数据的逻辑抽象),然后DStream被Spark Engine(Spark Core的离线计算引擎)执行并行处理,产生结果数据流DStream输出。
简言之,Spark Streaming就是先把数据按时间切分,然后按传统离线处理方式计算。
2)基本术语
(1)离散流(Discretized Stream)或DStream
这是Spark Streaming对内部持续的实时数据流的抽象描述,即处理的一个实时数据流,在Spark Streaming中对应于一个DStream实例。
(2)批数据(Batch Data)
将实时流数据以时间片为单位进行分批,将流处理转化为时间片数据的批处理。随着持续时间的推移,这些处理结果就形成了对应的结果数据流。
(3)时间片或批处理时间间隔(BatchInterval)
这是人为地对流数据进行定量的标准,以时间片作为拆分流数据的依据。一个时间片的数据组成一个DStream,对应一个RDD实例。
(4)窗口长度(Window Length)
窗口长度是一个窗口覆盖的流数据的时间长度。它必须是批处理时间间隔的倍数。
(5)滑动时间间隔(https://www.daowen.com)
滑动时间间隔是前一个窗口到后一个窗口所经过的时间长度。它必须是批处理时间间隔的倍数。
(6)Input DStream
一个Input DStream是一个特殊的DStream,它将Spark Streaming连接到一个外部数据源来读取数据。
3)Storm与Spark Streming比较
(1)处理模型以及延迟
虽然这两框架都提供了可扩展性(Scalability)和可容错性(Fault Tolerance),但是它们的处理模型从根本上来说是不一样的。Storm可以实现亚秒级的时延处理,而每次只处理一条事件(Event),而Spark Streaming却可以在一个短暂的时间窗口里面处理多条事件(Event)。所以说Storm可以实现亚秒级时延的处理,而Spark Streaming则是有一定的时延。
(2)容错和数据保证
Storm和Spark Streaming在容错数据保证上都作出了各自的权衡。但Spark Streaming在容错方面提供了对有状态的计算更好的支持。
在Storm中,每条记录在系统的移动过程中都需要被标记跟踪,所以Storm只能保证每条记录最少被处理一次,但是允许从错误状态恢复时被处理多次。这就意味着可变更的状态可能被更新两次从而导致结果不正确。Spark Streaming仅仅需要在批处理级别对记录进行追踪,所以它能保证每个批处理记录仅仅被处理一次,即使是node挂掉。尽管Storm的Trident Library可以保证一条记录被处理一次,但是它依赖于事务更新状态,而这个过程是很慢的,并且需要由用户去实现。
(3)实现和编程API
Storm主要是由Clojure语言实现的,Spark Streaming是由Scala实现的。Storm是由BackType和Twitter开发的,而Spark Streaming是由UC Berkeley开发的,它可以很好地与Spark批处理计算框架集成。
Storm提供了Java API,同时也支持其他语言的API。 Spark Streaming支持Scala和Java语言(其实也支持Python)。
(4)批处理框架集成
Spark Streaming有一个很棒的特性就是它是在Spark框架上运行的。这样就可以像使用其他批处理代码一样来写Spark Streaming程序,或者是在Spark中交互查询。减少了单独编写流批量处理程序和历史数据处理程序的工作量。
(5)生产支持
Storm已经出现好多年了,而且自从2011年开始就在Twitter内部生产环境中使用。而Spark Streaming是一个新的项目,并且在2013年仅被Sharethrough使用。
Storm是Hortonworks Hadoop数据平台中流处理的解决方案,而Spark Streaming出现在MapR的分布式平台和Cloudera的企业数据平台中。除此之外,Databricks是为Spark(包括Spark Streaming)提供技术支持的公司。
(6)集群管理集成
Storm和Spark Streaming都可以在各自的集群框架中运行,Storm也可以在Mesos上运行,而Spark Streaming也可以在YARN和Mesos上运行。