spark Dstreams-基本内容

AI 摘要 / TL;DR

spark Dstreams 基本内容 来源: https://www.ucloud.cn/yun/129640.html 作者: IT那活儿 发布日期: 发布于2023 01 11 13:20 spark Dstreams 基本内容 点击上方“IT那活儿”,关注后了解更多内容,不管IT什么活儿,干就完了!!! 简 介

来源: https://www.ucloud.cn/yun/129640.html 作者: IT那活儿 发布日期: 发布于2023-01-11 13:20

spark Dstreams-基本内容

点击上方“IT那活儿”,关注后了解更多内容,不管IT什么活儿,干就完了!!!

简  介

Spark Streaming是spark core的延伸,支持实时数据流的可伸缩、高吞吐量、容错机制,支持的数据源有kafka、Kinesis、TCP等),且可以使用复杂的算法(map,reduce,join,window等)处理数据,最后可以将处理后的数据推送到文件系统、数据库或仪表盘。(如下图,您也可以再数据流上应用spark的机器学习和图形算法。 )

01

在内部,spark streaming接受实时数据流并且将数据流切割成一个个的小批次,然后由spark引擎进行处理,以批量生成最终结果流。

Spark streaming提供一个高级抽象叫做离散化数据流Dstreams,他代表一个连续的数据流,可以从kafka,kinesis等源输入创建数据流,也可以通过对其他数据流应用高级操作来创建数据流。

在内部,一个Dstreams代表一系列的RDD。

02

具体步骤:

批处理间隔必须根据应用程序和可用集群资源进行设置,ssc也可以通过SparkContext创建:

val ssc = new StreamingContext(sc, Seconds(1))

创建ssc之后,按照下面步骤进行:

  1. 通过输入数据源创建Dstreams。
  2. 对数据流进行转换和输出操作。
  3. 使用ssc.start()启动应用程序。
  4. 使用ssc.awaitTermination()等待停止处理。(手动或由于任何错误)
  5. 使用ssc.stop()手动停止应用程序。

举例演示:

使用spark-streaming前提是添加依赖:

如果数据源为kafka,由于spark core中无kafka依赖,需要额外添加依赖:

StreamingContext是流功能的入口,我们创建一个本地StreamingContext,两个执行线程,批处理间隔为1秒。

val conf = new SparkConf().setMaster("local[2]").setAppName("Dstreams")

master是spark,Mesos,Kubernetes,或者YARN集群URL,或者本地模式运行的特殊local[*]字符串。

实际上,在集群上运行时,你不希望在程序中写死master,而是在使用spark-submit启动应用程序时在那里接收master。

03

离散化数据流(Dstreams)

Dstreams是spark streaming提供的基本抽象,它代表一个连续的数据流,不管是输入源创建的Dsteam还是通过转换生成的输入流,在内部Dstream代表一系列连续的RDD,每个RDD都是特定时间间隔的数据。

应用于数据流的任何操作都转化为底层RDD的操作,例如在之前的例子中,对用数据流中每个RDD应用flatMap操作,如图:

这些底层RDD转换由spark引擎计算,Dstream操作隐藏了大部分这些细节,并为开发人员提供了更高级别的API。

04

输入数据流(Input Dstreams)表示从源接受的输入数据流。

在上面的例子中,lines就是输入数据流(Input DStreams),代表从TCP socket接收的数据。

每个输入数据流(Input DStreams)除了文件流都与接收器对象关联,该对象从源接受的数据并将其在spark的内存中进行处理。

Spark streaming提供两类内置streaming源:

  1. 基本source:在StreamingContext中可以直接用的基本source,像file source,socket source。
  2. 高级source:像kafka,Kinesis等,需要添加额外依赖。

注意:

  1. 一旦StreamingContext启动,就不能设置或添加新的流计算。

  2. StreamingContext一旦停止就不能重启。

  3. 一个JVM中只能同时只能有一个StreamingContext,ssc.stop()也会停止SparkContext,如果仅停止streamingContext,需要ssc.stop(false)。

  4. 只要在创建下一个StreamingContext之前停止上一个StreamingContext而不停止SparkContext,就可以重新用于创建StreamingContetx。

  5. 如果希望在应用程序中并行接受多个数据流,则可以创建多个输入数据流,这将同时创建接受多个输入数据流的多个接收器,Spark worker/executor是一个长期运行的任务,因此它占用分配给Spark流应用程序的一个核心。

    因此,务必记住,Spark Streaming应用程序需要分配足够的内核(或线程,如果在本地运行)来处理接收到的数据,以及运行接收器。

  6. 在本地运行spark streaming应用程序时,不要将master设置成local[1]或local,这两种情况意味着只有一个线程用于本地运行任务,如果您使用的是基于接收器的输入数据流(例如套接字,kafka等),则使用单个线程运行接收器,就没有其他线程处理数据了。

  7. 在集群上运行时,分配给spark streaming应用程序的内核数必须大于接收器数,否则应用只能接受数据而无法处理数据。

本文作者:潘宗昊

本文来源:IT那活儿(上海新炬王翦团队)