4.3.5 Spark Streaming流计算编程案例

更新于 2026年10月10日 版权声明
4.3.5 Spark Streaming流计算编程案例

接下来以Spark Streaming官方提供的WordCount代码为例来介绍Spark Streaming的使用方式,实现代码如图4⁃21所示。

可以看出,代码非常简单,程序流程也非常清晰,这正好体现了Spark Streaming具有简洁易用的编程优势,下面简单介绍一下该程序的关键部分。

图示

图4⁃21 Spark Streaming使用方式

1)创建StreamingContext对象

同Spark初始化需要创建SparkContext对象一样,使用Spark Streaming就需要创建StreamingContext对象。创建StreamingContext对象所需的参数与SparkContext基本一致,包括指明Master,设定名称(如NetworkWordCount)。需要注意的是参数Seconds(1),Spark Streaming需要指定处理数据的时间间隔,如上例所示的1 s,那么Spark Streaming会以1 s为时间窗口进行数据处理。此参数需要根据用户的需求和集群的处理能力进行适当的设置。

2)创建InputDStream(https://www.daowen.com)

如同Storm的Spout,Spark Streaming需要指明数据源。如上例的socketTextStream,Spark Streaming以Socket连接作为数据源读取数据。当然Spark Streaming支持多种不同的数据源,包括Kafka、 Flume、HDFS/S3、Kinesis和Twitter等数据源。

3)操作DStream

对于从数据源得到的DStream,用户可以在其基础上进行各种操作,如上例的操作就是一个典型的WordCount执行流程。对于当前时间窗口内从数据源得到的数据首先进行分割,然后利用map和reduceByKey方法进行计算,最后使用print方法输出结果。

4)启动Spark Streaming

之前所作的所有步骤只是创建了执行流程,程序没有真正连接上数据源,也没有对数据进行任何操作,只是设定好了所有的执行计划,当ssc.start()启动后,程序才真正进行所有预期的操作。

↑上一章 ↓下一章
关注公众号获取验证码
复制内容需要验证码(7.99元/天)