Spark Streaming 编程指南
注意
Spark Streaming 是 Spark 上一代的流处理引擎。目前 Spark Streaming 已不再更新,属于遗留项目。Spark 中有一个更新、更易于使用的流处理引擎,称为 Structured Streaming。你应该在流处理应用程序和管道中使用 Spark Structured Streaming。请参阅 Structured Streaming 编程指南。
概述
Spark Streaming 是核心 Spark API 的扩展,支持对实时数据流进行可扩展、高吞吐、容错的流处理。数据可以通过 Kafka、Kinesis 或 TCP 套接字等多种来源摄入,并可以使用类似 map、reduce、join 和 window 等高级函数表示的复杂算法进行处理。最后,处理后的数据可以输出到文件系统、数据库和实时仪表板中。实际上,你可以在数据流上应用 Spark 的机器学习和图处理算法。
从内部实现来看,其工作原理如下:Spark Streaming 接收实时输入数据流,并将数据划分为多个批次,然后由 Spark 引擎处理这些批次,最终生成批处理形式的结果流。
Spark Streaming 提供了一种称为离散流或DStream的高级抽象,它代表一个连续的数据流。DStream 可以通过 Kafka 和 Kinesis 等源的输入数据流创建,也可以通过对其他 DStream 应用高级操作来创建。在内部,DStream 被表示为一系列 RDDs。
本指南将向你展示如何开始使用 DStream 编写 Spark Streaming 程序。你可以使用 Scala、Java 或 Python(在 Spark 1.2 中引入)编写 Spark Streaming 程序,本指南中均有介绍。本指南中会包含选项卡,让你可以在不同语言的代码片段之间进行切换。
注意: 有少数 API 在 Python 中要么不同,要么不可用。在本指南中,你会看到标记 Python API 用以突出显示这些差异。
快速入门示例
在深入了解如何编写自己的 Spark Streaming 程序之前,让我们先快速看一下一个简单的 Spark Streaming 程序是什么样的。假设我们想统计从监听 TCP 套接字的数据服务器接收到的文本数据中的单词数量。你只需要执行以下操作。
首先,我们导入 StreamingContext,它是所有流处理功能的主要入口点。我们创建一个包含两个执行线程的本地 StreamingContext,批处理间隔为 1 秒。
from pyspark import SparkContext
from pyspark.streaming import StreamingContext
# Create a local StreamingContext with two working thread and batch interval of 1 second
sc = SparkContext("local[2]", "NetworkWordCount")
ssc = StreamingContext(sc, 1)使用此上下文,我们可以创建一个表示来自 TCP 源(指定主机名,例如 localhost,和端口,例如 9999)的流数据的 DStream。
# Create a DStream that will connect to hostname:port, like localhost:9999
lines = ssc.socketTextStream("localhost", 9999)这个 lines DStream 代表将从数据服务器接收到的数据流。此 DStream 中的每条记录都是一行文本。接下来,我们要按空格将行拆分为单词。
# Split each line into words
words = lines.flatMap(lambda line: line.split(" "))flatMap 是一种一对多的 DStream 操作,它通过从源 DStream 中的每个记录生成多个新记录来创建一个新的 DStream。在本例中,每一行将被拆分为多个单词,单词流由 words DStream 表示。接下来,我们要统计这些单词。
# Count each word in each batch
pairs = words.map(lambda word: (word, 1))
wordCounts = pairs.reduceByKey(lambda x, y: x + y)
# Print the first ten elements of each RDD generated in this DStream to the console
wordCounts.pprint()words DStream 进一步被映射(一对一转换)为 (word, 1) 对的 DStream,然后进行归约以获取每个数据批次中单词的频率。最后,wordCounts.pprint() 将打印每秒生成的计数中的一部分。
注意,当执行这些代码行时,Spark Streaming 只是设置了它在启动后将执行的计算,此时并没有开始真正的处理。在完成所有转换设置后,要开始处理,我们最后调用:
ssc.start() # Start the computation
ssc.awaitTermination() # Wait for the computation to terminate完整代码可以在 Spark Streaming 示例 NetworkWordCount 中找到。
首先,我们将 Spark Streaming 类的名称和一些来自 StreamingContext 的隐式转换导入到我们的环境中,以便为我们需要的其他类(如 DStream)添加有用的方法。StreamingContext 是所有流处理功能的主要入口点。我们创建一个包含两个执行线程的本地 StreamingContext,批处理间隔为 1 秒。
import org.apache.spark._
import org.apache.spark.streaming._
import org.apache.spark.streaming.StreamingContext._ // not necessary since Spark 1.3
// Create a local StreamingContext with two working thread and batch interval of 1 second.
// The master requires 2 cores to prevent a starvation scenario.
val conf = new SparkConf().setMaster("local[2]").setAppName("NetworkWordCount")
val ssc = new StreamingContext(conf, Seconds(1))使用此上下文,我们可以创建一个表示来自 TCP 源(指定主机名,例如 localhost,和端口,例如 9999)的流数据的 DStream。
// Create a DStream that will connect to hostname:port, like localhost:9999
val lines = ssc.socketTextStream("localhost", 9999)这个 lines DStream 代表将从数据服务器接收到的数据流。此 DStream 中的每条记录都是一行文本。接下来,我们要按空格字符将行拆分为单词。
// Split each line into words
val words = lines.flatMap(_.split(" "))flatMap 是一种一对多的 DStream 操作,它通过从源 DStream 中的每个记录生成多个新记录来创建一个新的 DStream。在本例中,每一行将被拆分为多个单词,单词流由 words DStream 表示。接下来,我们要统计这些单词。
import org.apache.spark.streaming.StreamingContext._ // not necessary since Spark 1.3
// Count each word in each batch
val pairs = words.map(word => (word, 1))
val wordCounts = pairs.reduceByKey(_ + _)
// Print the first ten elements of each RDD generated in this DStream to the console
wordCounts.print()words DStream 进一步被映射(一对一转换)为 (word, 1) 对的 DStream,然后进行归约以获取每个数据批次中单词的频率。最后,wordCounts.print() 将打印每秒生成的计数中的一部分。
注意,当执行这些代码行时,Spark Streaming 只是设置了它在启动后将执行的计算,此时并没有开始真正的处理。在完成所有转换设置后,要开始处理,我们最后调用:
ssc.start() // Start the computation
ssc.awaitTermination() // Wait for the computation to terminate完整代码可以在 Spark Streaming 示例 NetworkWordCount 中找到。
首先,我们创建一个 JavaStreamingContext 对象,它是所有流处理功能的主要入口点。我们创建一个包含两个执行线程的本地 StreamingContext,批处理间隔为 1 秒。
import org.apache.spark.*;
import org.apache.spark.api.java.function.*;
import org.apache.spark.streaming.*;
import org.apache.spark.streaming.api.java.*;
import scala.Tuple2;
// Create a local StreamingContext with two working thread and batch interval of 1 second
SparkConf conf = new SparkConf().setMaster("local[2]").setAppName("NetworkWordCount");
JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(1));使用此上下文,我们可以创建一个表示来自 TCP 源(指定主机名,例如 localhost,和端口,例如 9999)的流数据的 DStream。
// Create a DStream that will connect to hostname:port, like localhost:9999
JavaReceiverInputDStream<String> lines = jssc.socketTextStream("localhost", 9999);这个 lines DStream 代表将从数据服务器接收到的数据流。此流中的每条记录都是一行文本。然后,我们要按空格将行拆分为单词。
// Split each line into words
JavaDStream<String> words = lines.flatMap(x -> Arrays.asList(x.split(" ")).iterator());flatMap 是一种 DStream 操作,它通过从源 DStream 中的每个记录生成多个新记录来创建一个新的 DStream。在本例中,每一行将被拆分为多个单词,单词流由 words DStream 表示。注意,我们使用 FlatMapFunction 对象定义了转换。正如我们稍后会发现的那样,Java API 中有许多此类便捷类,有助于定义 DStream 转换。
接下来,我们要统计这些单词。
// Count each word in each batch
JavaPairDStream<String, Integer> pairs = words.mapToPair(s -> new Tuple2<>(s, 1));
JavaPairDStream<String, Integer> wordCounts = pairs.reduceByKey((i1, i2) -> i1 + i2);
// Print the first ten elements of each RDD generated in this DStream to the console
wordCounts.print();words DStream 使用 PairFunction 对象进一步映射(一对一转换)为 (word, 1) 对的 DStream。然后,使用 Function2 对象将其归约以获取每个数据批次中单词的频率。最后,wordCounts.print() 将打印每秒生成的计数中的一部分。
注意,当执行这些代码行时,Spark Streaming 只是设置了它在启动后将执行的计算,此时并没有开始真正的处理。在完成所有转换设置后,要开始处理,我们最后调用 start 方法。
jssc.start(); // Start the computation
jssc.awaitTermination(); // Wait for the computation to terminate完整代码可以在 Spark Streaming 示例 JavaNetworkWordCount 中找到。
如果你已经下载并构建了 Spark,你可以按照以下方式运行此示例。你首先需要通过以下方式运行 Netcat(大多数类 Unix 系统中都有的一个小工具)作为数据服务器:
$ nc -lk 9999然后,在另一个终端中,你可以通过以下方式启动该示例:
$ ./bin/spark-submit examples/src/main/python/streaming/network_wordcount.py localhost 9999$ ./bin/run-example streaming.NetworkWordCount localhost 9999$ ./bin/run-example streaming.JavaNetworkWordCount localhost 9999之后,在运行 netcat 服务器的终端中输入的任何行都将被统计并每秒打印在屏幕上。它看起来会像下面这样。
|
|
基本概念
接下来,我们将超越简单的示例,详细说明 Spark Streaming 的基础知识。
链接依赖
与 Spark 类似,Spark Streaming 可通过 Maven Central 获取。要编写自己的 Spark Streaming 程序,你需要将以下依赖项添加到你的 SBT 或 Maven 项目中。
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming_2.13</artifactId>
<version>4.1.1</version>
<scope>provided</scope>
</dependency>
libraryDependencies += "org.apache.spark" % "spark-streaming_2.13" % "4.1.1" % "provided"
对于从 Spark Streaming 核心 API 中未提供的 Kafka 和 Kinesis 等源摄入数据,你需要将相应的构件 spark-streaming-xyz_2.13 添加到依赖项中。例如,一些常见的如下所示。
| 源 | 构件 (Artifact) |
|---|---|
| Kafka | spark-streaming-kafka-0-10_2.13 |
| Kinesis | spark-streaming-kinesis-asl_2.13 [Amazon 软件许可] |
有关最新列表,请参阅 Maven 存储库,以获取支持的源和构件的完整列表。
初始化 StreamingContext
要初始化 Spark Streaming 程序,必须创建一个 StreamingContext 对象,它是所有 Spark Streaming 功能的主要入口点。
StreamingContext 对象可以从 SparkContext 对象创建。
from pyspark import SparkContext
from pyspark.streaming import StreamingContext
sc = SparkContext(master, appName)
ssc = StreamingContext(sc, 1)appName 参数是你应用程序的名称,显示在集群 UI 上。master 是一个 Spark 或 YARN 集群 URL,或者是在本地模式下运行的特殊 “local[*]” 字符串。在实践中,当在集群上运行时,你不会想在程序中硬编码 master,而是使用 spark-submit 启动应用程序并在那里接收它。然而,对于本地测试和单元测试,你可以传递“local[*]”以进程内方式运行 Spark Streaming(检测本地系统中的核心数量)。
批处理间隔必须根据应用程序的延迟要求和可用的集群资源来设置。有关更多详细信息,请参阅 性能调优 部分。
StreamingContext 对象可以从 SparkConf 对象创建。
import org.apache.spark._
import org.apache.spark.streaming._
val conf = new SparkConf().setAppName(appName).setMaster(master)
val ssc = new StreamingContext(conf, Seconds(1))appName 参数是你应用程序的名称,显示在集群 UI 上。master 是一个 Spark、Kubernetes 或 YARN 集群 URL,或者是在本地模式下运行的特殊 “local[*]” 字符串。在实践中,当在集群上运行时,你不会想在程序中硬编码 master,而是使用 spark-submit 启动应用程序并在那里接收它。然而,对于本地测试和单元测试,你可以传递“local[*]”以进程内方式运行 Spark Streaming(检测本地系统中的核心数量)。注意,这在内部创建了一个 SparkContext(所有 Spark 功能的起点),可以通过 ssc.sparkContext 访问它。
批处理间隔必须根据应用程序的延迟要求和可用的集群资源来设置。有关更多详细信息,请参阅 性能调优 部分。
StreamingContext 对象也可以从现有的 SparkContext 对象创建。
import org.apache.spark.streaming._
val sc = ... // existing SparkContext
val ssc = new StreamingContext(sc, Seconds(1))JavaStreamingContext 对象可以从 SparkConf 对象创建。
import org.apache.spark.*;
import org.apache.spark.streaming.api.java.*;
SparkConf conf = new SparkConf().setAppName(appName).setMaster(master);
JavaStreamingContext ssc = new JavaStreamingContext(conf, new Duration(1000));appName 参数是你应用程序的名称,显示在集群 UI 上。master 是一个 Spark 或 YARN 集群 URL,或者是在本地模式下运行的特殊 “local[*]” 字符串。在实践中,当在集群上运行时,你不会想在程序中硬编码 master,而是使用 spark-submit 启动应用程序并在那里接收它。然而,对于本地测试和单元测试,你可以传递“local[*]”以进程内方式运行 Spark Streaming。注意,这在内部创建了一个 JavaSparkContext(所有 Spark 功能的起点),可以通过 ssc.sparkContext 访问它。
批处理间隔必须根据应用程序的延迟要求和可用的集群资源来设置。有关更多详细信息,请参阅 性能调优 部分。
JavaStreamingContext 对象也可以从现有的 JavaSparkContext 创建。
import org.apache.spark.streaming.api.java.*;
JavaSparkContext sc = ... //existing JavaSparkContext
JavaStreamingContext ssc = new JavaStreamingContext(sc, Durations.seconds(1));定义上下文后,你必须执行以下操作:
- 通过创建输入 DStreams 来定义输入源。
- 通过对 DStreams 应用转换和输出操作来定义流处理计算。
- 使用
streamingContext.start()开始接收数据并对其进行处理。 - 使用
streamingContext.awaitTermination()等待处理停止(手动或由于任何错误)。 - 可以使用
streamingContext.stop()手动停止处理。
需要记住的要点
- 一旦上下文启动,就不能对其设置或添加新的流计算。
- 一旦上下文停止,它就无法重新启动。
- JVM 中同一时间只能有一个 StreamingContext 处于活动状态。
- 对 StreamingContext 调用
stop()也会停止 SparkContext。如果只想停止 StreamingContext,请将stop()的可选参数stopSparkContext设置为 false。 - 只要在创建下一个 StreamingContext 之前先停止上一个 StreamingContext(不停止 SparkContext),就可以重复使用 SparkContext 来创建多个 StreamingContext。
离散流 (DStreams)
离散流 (Discretized Stream) 或 DStream 是 Spark Streaming 提供的基本抽象。它代表一个连续的数据流,可以是接收到的输入数据流,也可以是通过转换输入流生成的已处理数据流。在内部,DStream 由连续的 RDD 系列表示,RDD 是 Spark 的不可变分布式数据集抽象(更多详细信息请参阅 Spark 编程指南)。DStream 中的每个 RDD 都包含来自特定间隔的数据,如下图所示。
应用于 DStream 的任何操作都会转换为对底层 RDD 的操作。例如,在前面的示例中,将行流转换为单词,flatMap 操作应用于 lines DStream 中的每个 RDD,以生成 words DStream 的 RDD。如下图所示。
这些底层的 RDD 转换由 Spark 引擎计算。DStream 操作隐藏了大部分细节,为开发人员提供了一个更高级的 API 以方便使用。这些操作将在后续章节中详细讨论。
输入 DStreams 和接收器 (Receivers)
输入 DStreams 是表示从流式源接收的输入数据流的 DStreams。在快速示例中,lines 是一个输入 DStream,因为它表示从 netcat 服务器接收到的数据流。每个输入 DStream(文件流除外,本节后面会讨论)都与一个 Receiver(Scala 文档,Java 文档)对象关联,该对象从源接收数据并将其存储在 Spark 的内存中以进行处理。
Spark Streaming 提供两类内置流式源。
- 基础源:直接在 StreamingContext API 中可用的源。例如:文件系统和套接字连接。
- 高级源:Kafka、Kinesis 等源通过额外的实用程序类提供。这些需要链接额外的依赖项,如链接依赖一节中所述。
我们将在本节后面讨论每个类别中存在的一些源。
注意,如果你希望在流式应用程序中并行接收多个数据流,可以创建多个输入 DStreams(在性能调优一节中进一步讨论)。这将创建多个接收器,同时接收多个数据流。但请注意,Spark worker/executor 是一个长期运行的任务,因此它会占用分配给 Spark Streaming 应用程序的一个核心。因此,务必记住,Spark Streaming 应用程序需要分配足够的核心(如果本地运行则是线程)来处理接收到的数据,并运行接收器。
需要记住的要点
-
在本地运行 Spark Streaming 程序时,不要使用“local”或“local[1]”作为主 URL。这两者都意味着本地运行任务只会使用一个线程。如果你使用的是基于接收器的输入 DStream(例如套接字、Kafka 等),那么该单线程将被用于运行接收器,从而没有线程可用于处理接收到的数据。因此,在本地运行时,务必使用“local[n]”作为主 URL,其中 n > 要运行的接收器数量(有关如何设置 master 的信息,请参阅 Spark 属性)。
-
将逻辑扩展到集群上运行时,分配给 Spark Streaming 应用程序的核心数必须超过接收器的数量。否则系统将能够接收数据,但无法对其进行处理。
基础源
我们已经在快速示例中了解了 ssc.socketTextStream(...),它从通过 TCP 套接字连接接收到的文本数据创建一个 DStream。除了套接字之外,StreamingContext API 还提供了从文件作为输入源创建 DStream 的方法。
文件流 (File Streams)
对于从与 HDFS API 兼容的任何文件系统(即 HDFS、S3、NFS 等)上的文件读取数据,可以通过 StreamingContext.fileStream[KeyClass, ValueClass, InputFormatClass] 创建 DStream。
文件流不需要运行接收器,因此无需为接收文件数据分配任何核心。
对于简单的文本文件,最简单的方法是 StreamingContext.textFileStream(dataDirectory)。
fileStream 在 Python API 中不可用;只有 textFileStream 可用。
streamingContext.textFileStream(dataDirectory)streamingContext.fileStream[KeyClass, ValueClass, InputFormatClass](dataDirectory)对于文本文件
streamingContext.textFileStream(dataDirectory)streamingContext.fileStream<KeyClass, ValueClass, InputFormatClass>(dataDirectory);对于文本文件
streamingContext.textFileStream(dataDirectory);目录是如何被监控的
Spark Streaming 将监控目录 dataDirectory 并处理在该目录中创建的所有文件。
- 可以监控一个简单的目录,例如
"hdfs://namenode:8040/logs/"。在此类路径下直接发现的所有文件都将被处理。 - 可以提供一个 POSIX glob 模式,例如
"hdfs://namenode:8040/logs/2017/*"。这里,DStream 将由匹配该模式的目录中的所有文件组成。也就是说:它是目录的模式,而不是目录中文件的模式。 - 所有文件必须采用相同的数据格式。
- 文件被视为时间段的一部分是基于其修改时间,而不是创建时间。
- 处理后,当前窗口内对文件的更改不会导致文件被重新读取。也就是说:更新会被忽略。
- 目录下的文件越多,扫描更改所需的时间就越长——即使没有文件被修改。
- 如果使用通配符来识别目录,例如
"hdfs://namenode:8040/logs/2016-*",则将整个目录重命名以匹配该路径会将该目录添加到监控目录列表中。只有修改时间在当前窗口内的目录中的文件才会被包含在流中。 - 调用
FileSystem.setTimes()来固定时间戳是一种在稍后的窗口中拾取文件的方法,即使其内容没有改变。
使用对象存储作为数据源
“完整”文件系统(如 HDFS)往往在输出流创建后立即设置其文件的修改时间。当一个文件被打开时,即使在数据完全写入之前,它也可能被包含在 DStream 中——在此之后,同一窗口内对文件的更新将被忽略。也就是说:更改可能会被遗漏,数据会被从流中省略。
为了保证更改在窗口内被拾取,请将文件写入未监控的目录,然后,在输出流关闭后立即将其重命名到目标目录。只要重命名的文件在创建窗口期间出现在扫描的目标目录中,新数据就会被拾取。
相反,对象存储(如 Amazon S3 和 Azure Storage)通常具有缓慢的重命名操作,因为数据实际上是被复制的。此外,重命名后的对象可能将其 rename() 操作的时间作为其修改时间,因此可能不被视为原始创建时间所暗示的窗口的一部分。
需要针对目标对象存储进行仔细测试,以验证存储的时间戳行为是否与 Spark Streaming 预期的一致。对于通过所选对象存储进行流式传输的数据,直接写入目标目录可能是一种合适的策略。
有关此主题的更多详细信息,请参阅 Hadoop 文件系统规范。
基于自定义接收器的流
可以使用通过自定义接收器接收的数据流创建 DStream。有关详细信息,请参阅自定义接收器指南。
作为流的 RDD 队列
为了使用测试数据测试 Spark Streaming 应用程序,还可以使用 streamingContext.queueStream(queueOfRDDs) 创建基于 RDD 队列的 DStream。推入队列的每个 RDD 都将被视为 DStream 中的一个数据批次,并像流一样进行处理。
有关套接字和文件流的更多详细信息,请参阅相关函数的 API 文档:Python 中的 StreamingContext,Scala 中的 StreamingContext,以及 Java 中的 JavaStreamingContext。
高级源
Python API 截至 Spark 4.1.1,在这些源中,Kafka 和 Kinesis 在 Python API 中可用。
此类源需要与外部非 Spark 库进行接口,其中一些具有复杂的依赖关系(例如 Kafka)。因此,为了最大限度地减少与依赖版本冲突相关的问题,从这些源创建 DStream 的功能已被移动到单独的库中,可以在必要时显式地链接这些库。
注意,这些高级源在 Spark shell 中不可用,因此基于这些高级源的应用程序无法在 shell 中进行测试。如果你确实想在 Spark shell 中使用它们,则必须下载相应的 Maven 构件的 JAR 及其依赖项,并将它们添加到类路径中。
其中一些高级源如下所示。
-
Kafka: Spark Streaming 4.1.1 与 Kafka broker 版本 0.10 或更高版本兼容。有关详细信息,请参阅 Kafka 集成指南。
-
Kinesis: Spark Streaming 4.1.1 与 Kinesis Client Library 1.2.1 兼容。有关详细信息,请参阅 Kinesis 集成指南。
自定义源
Python API Python 中尚不支持此功能。
输入 DStreams 也可以由自定义数据源创建。你所要做的就是实现一个用户定义的 receiver(请参阅下一节以了解其含义),它可以从自定义源接收数据并将其推送到 Spark。有关详细信息,请参阅 自定义接收器指南。
接收器可靠性
根据数据源的 可靠性,可以有两种类型的数据源。源(如 Kafka)允许确认传输的数据。如果接收来自这些 可靠 源数据的系统能正确确认接收到的数据,则可以确保不会因为任何类型的故障而丢失任何数据。这导致了两种类型的接收器:
- 可靠接收器 - 当数据已被接收并存储在 Spark 中并进行复制后,可靠接收器会向可靠源正确发送确认。
- 不可靠接收器 - 不可靠接收器 不向源发送确认。这可用于不支持确认的源,甚至可用于在不想或不需要处理确认复杂性的情况下的可靠源。
有关如何编写可靠接收器的详细信息,请参见 自定义接收器指南。
DStream 上的转换操作
与 RDD 类似,转换允许修改输入 DStream 中的数据。DStreams 支持常规 Spark RDD 上提供的许多转换。一些常见的转换如下所示。
| 转换 | 含义 |
|---|---|
| map(func) | 通过将源 DStream 的每个元素传递给函数 func 来返回一个新的 DStream。 |
| flatMap(func) | 与 map 类似,但每个输入项可以映射到 0 个或多个输出项。 |
| filter(func) | 通过仅选择源 DStream 中 func 返回 true 的记录来返回一个新的 DStream。 |
| repartition(numPartitions) | 通过创建更多或更少的分区来更改此 DStream 中的并行度。 |
| union(otherStream) | 返回一个新的 DStream,其中包含源 DStream 和 otherDStream 中元素的并集。 |
| count() | 通过计算源 DStream 中每个 RDD 的元素数量,返回一个新的单元素 RDD 的 DStream。 |
| reduce(func) | 通过使用函数 func(接受两个参数并返回一个)聚合源 DStream 中每个 RDD 中的元素,返回一个新的单元素 RDD 的 DStream。该函数应该是结合律和交换律的,以便可以并行计算。 |
| countByValue() | 当在 K 类型元素的 DStream 上调用时,返回一个新的 (K, Long) 对的 DStream,其中每个键的值是其在源 DStream 的每个 RDD 中的频率。 |
| reduceByKey(func, [numTasks]) | 当在 (K, V) 对的 DStream 上调用时,返回一个新的 (K, V) 对的 DStream,其中每个键的值使用给定的 reduce 函数进行聚合。 注意: 默认情况下,这使用 Spark 的默认并行任务数(本地模式为 2,集群模式下数量由配置属性 spark.default.parallelism 决定)来进行分组。你可以传递可选的 numTasks 参数来设置不同的任务数。 |
| join(otherStream, [numTasks]) | 当在 (K, V) 和 (K, W) 对的两个 DStreams 上调用时,返回一个新的 (K, (V, W)) 对的 DStream,其中包含每个键的所有元素对。 |
| cogroup(otherStream, [numTasks]) | 当在 (K, V) 和 (K, W) 对的 DStream 上调用时,返回一个新的 (K, Seq[V], Seq[W]) 元组的 DStream。 |
| transform(func) | 通过对源 DStream 的每个 RDD 应用 RDD 到 RDD 的函数来返回一个新的 DStream。这可用于在 DStream 上执行任意 RDD 操作。 |
| updateStateByKey(func) | 返回一个新的“状态” DStream,其中每个键的状态是通过将给定函数应用于键的先前状态和键的新值来更新的。这可用于为每个键维护任意状态数据。 |
其中一些转换值得更详细地讨论。
UpdateStateByKey 操作
updateStateByKey 操作允许你维护任意状态,同时用新信息不断更新它。要使用此功能,你必须执行两个步骤:
- 定义状态 - 状态可以是任何数据类型。
- 定义状态更新函数 - 使用函数指定如何利用先前状态和来自输入流的新值来更新状态。
在每个批次中,Spark 将为所有现有键应用状态更新函数,无论它们在批次中是否有新数据。如果更新函数返回 None,则键值对将被消除。
让我们用一个例子来说明这一点。假设你想维护文本数据流中每个单词的运行计数。在这里,运行计数就是状态,它是一个整数。我们定义更新函数如下:
def updateFunction(newValues, runningCount):
if runningCount is None:
runningCount = 0
return sum(newValues, runningCount) # add the new values with the previous running count to get the new count这应用于包含单词的 DStream(例如,前面的示例中包含 (word, 1) 对的 pairs DStream)。
runningCounts = pairs.updateStateByKey(updateFunction)更新函数将为每个单词调用,其中 newValues 具有一系列 1(来自 (word, 1) 对),而 runningCount 具有之前的计数。对于完整的 Python 代码,请查看示例 stateful_network_wordcount.py。
def updateFunction(newValues: Seq[Int], runningCount: Option[Int]): Option[Int] = {
val newCount = ... // add the new values with the previous running count to get the new count
Some(newCount)
}这应用于包含单词的 DStream(例如,前面的示例中包含 (word, 1) 对的 pairs DStream)。
val runningCounts = pairs.updateStateByKey[Int](updateFunction _)更新函数将为每个单词调用,其中 newValues 具有一系列 1(来自 (word, 1) 对),而 runningCount 具有之前的计数。
Function2<List<Integer>, Optional<Integer>, Optional<Integer>> updateFunction =
(values, state) -> {
Integer newSum = ... // add the new values with the previous running count to get the new count
return Optional.of(newSum);
};这应用于包含单词的 DStream(例如,快速示例中包含 (word, 1) 对的 pairs DStream)。
JavaPairDStream<String, Integer> runningCounts = pairs.updateStateByKey(updateFunction);更新函数将为每个单词调用,其中 newValues 具有一系列 1(来自 (word, 1) 对),而 runningCount 具有之前的计数。对于完整的 Java 代码,请查看示例 JavaStatefulNetworkWordCount.java。
注意,使用 updateStateByKey 需要配置检查点目录,这将在 检查点 部分详细讨论。
Transform 操作
transform 操作(连同其变体如 transformWith)允许将任意 RDD 到 RDD 的函数应用于 DStream。它可用于应用 DStream API 中未公开的任何 RDD 操作。例如,将数据流中的每个批次与另一个数据集进行 join 的功能在 DStream API 中没有直接公开。但是,你可以轻松地使用 transform 来实现这一点。这开启了非常强大的可能性。例如,可以通过将输入数据流与预先计算的垃圾邮件信息(可能也是用 Spark 生成的)进行连接,然后基于它进行过滤来执行实时数据清洗。
spamInfoRDD = sc.pickleFile(...) # RDD containing spam information
# join data stream with spam information to do data cleaning
cleanedDStream = wordCounts.transform(lambda rdd: rdd.join(spamInfoRDD).filter(...))val spamInfoRDD = ssc.sparkContext.newAPIHadoopRDD(...) // RDD containing spam information
val cleanedDStream = wordCounts.transform { rdd =>
rdd.join(spamInfoRDD).filter(...) // join data stream with spam information to do data cleaning
...
}import org.apache.spark.streaming.api.java.*;
// RDD containing spam information
JavaPairRDD<String, Double> spamInfoRDD = jssc.sparkContext().newAPIHadoopRDD(...);
JavaPairDStream<String, Integer> cleanedDStream = wordCounts.transform(rdd -> {
rdd.join(spamInfoRDD).filter(...); // join data stream with spam information to do data cleaning
...
});注意,提供的函数在每个批次间隔中都会被调用。这允许你执行随时间变化的 RDD 操作,也就是说,RDD 操作、分区数量、广播变量等可以在批次之间更改。
窗口操作
Spark Streaming 还提供了 窗口计算,允许你在滑动数据窗口上应用转换。下图说明了这种滑动窗口。
如图所示,每次窗口在源 DStream 上 滑动 时,落在窗口内的源 RDD 都会被组合并运行以生成窗口化 DStream 的 RDD。在这种特定情况下,该操作应用于最后 3 个时间单位的数据,并按 2 个时间单位滑动。这表明任何窗口操作都需要指定两个参数:
- 窗口长度 - 窗口的持续时间(图中为 3)。
- 滑动间隔 - 执行窗口操作的间隔(图中为 2)。
这两个参数必须是源 DStream 批处理间隔(图中为 1)的倍数。
让我们用一个例子来说明窗口操作。假设你想通过每 10 秒生成过去 30 秒数据的单词计数来扩展前面的示例。为此,我们必须在过去 30 秒数据的 (word, 1) 对的 pairs DStream 上应用 reduceByKey 操作。这是使用操作 reduceByKeyAndWindow 完成的。
# Reduce last 30 seconds of data, every 10 seconds
windowedWordCounts = pairs.reduceByKeyAndWindow(lambda x, y: x + y, lambda x, y: x - y, 30, 10)// Reduce last 30 seconds of data, every 10 seconds
val windowedWordCounts = pairs.reduceByKeyAndWindow((a:Int,b:Int) => (a + b), Seconds(30), Seconds(10))// Reduce last 30 seconds of data, every 10 seconds
JavaPairDStream<String, Integer> windowedWordCounts = pairs.reduceByKeyAndWindow((i1, i2) -> i1 + i2, Durations.seconds(30), Durations.seconds(10));一些常见的窗口操作如下所示。所有这些操作都采用所述的两个参数 - windowLength 和 slideInterval。
| 转换 | 含义 |
|---|---|
| window(windowLength, slideInterval) | 返回一个基于源 DStream 的窗口化批次计算的新 DStream。 |
| countByWindow(windowLength, slideInterval) | 返回流中元素的滑动窗口计数。 |
| reduceByWindow(func, windowLength, slideInterval) | 返回一个新的单元素流,通过使用 func 在滑动间隔内聚合流中的元素来创建。该函数应该是结合律和交换律的,以便可以并行正确计算。 |
| reduceByKeyAndWindow(func, windowLength, slideInterval, [numTasks]) | 当在 (K, V) 对的 DStream 上调用时,返回一个新的 (K, V) 对的 DStream,其中每个键的值使用给定的 reduce 函数 func 在滑动窗口中的批次上进行聚合。 注意: 默认情况下,这使用 Spark 的默认并行任务数(本地模式为 2,集群模式下数量由配置属性 spark.default.parallelism 决定)来进行分组。你可以传递可选的 numTasks 参数来设置不同的任务数。 |
| reduceByKeyAndWindow(func, invFunc, windowLength, slideInterval, [numTasks]) |
上述 |
| countByValueAndWindow(windowLength, slideInterval, [numTasks]) | 当在 (K, V) 对的 DStream 上调用时,返回一个新的 (K, Long) 对的 DStream,其中每个键的值是其在滑动窗口内的频率。与 reduceByKeyAndWindow 一样,reduce 任务的数量可以通过可选参数进行配置。 |
Join 操作
最后,值得强调的是你可以多么轻松地在 Spark Streaming 中执行不同类型的连接。
流-流 Join
流可以非常容易地与其他流连接。
stream1 = ...
stream2 = ...
joinedStream = stream1.join(stream2)val stream1: DStream[String, String] = ...
val stream2: DStream[String, String] = ...
val joinedStream = stream1.join(stream2)JavaPairDStream<String, String> stream1 = ...
JavaPairDStream<String, String> stream2 = ...
JavaPairDStream<String, Tuple2<String, String>> joinedStream = stream1.join(stream2);在这里,在每个批次间隔中,stream1 生成的 RDD 将与 stream2 生成的 RDD 连接。你也可以执行 leftOuterJoin、rightOuterJoin、fullOuterJoin。此外,在流的窗口上进行连接通常非常有用。这也相当容易。
windowedStream1 = stream1.window(20)
windowedStream2 = stream2.window(60)
joinedStream = windowedStream1.join(windowedStream2)val windowedStream1 = stream1.window(Seconds(20))
val windowedStream2 = stream2.window(Minutes(1))
val joinedStream = windowedStream1.join(windowedStream2)JavaPairDStream<String, String> windowedStream1 = stream1.window(Durations.seconds(20));
JavaPairDStream<String, String> windowedStream2 = stream2.window(Durations.minutes(1));
JavaPairDStream<String, Tuple2<String, String>> joinedStream = windowedStream1.join(windowedStream2);流-数据集 Join
这在解释 DStream.transform 操作时已经介绍过。这是将窗口化流与数据集进行连接的另一个示例。
dataset = ... # some RDD
windowedStream = stream.window(20)
joinedStream = windowedStream.transform(lambda rdd: rdd.join(dataset))val dataset: RDD[String, String] = ...
val windowedStream = stream.window(Seconds(20))...
val joinedStream = windowedStream.transform { rdd => rdd.join(dataset) }JavaPairRDD<String, String> dataset = ...
JavaPairDStream<String, String> windowedStream = stream.window(Durations.seconds(20));
JavaPairDStream<String, String> joinedStream = windowedStream.transform(rdd -> rdd.join(dataset));实际上,你也可以动态更改要连接的数据集。传递给 transform 的函数在每个批次间隔都会被评估,因此将使用 dataset 引用当前指向的数据集。
DStream 转换的完整列表可在 API 文档中找到。有关 Python API,请参阅 DStream。有关 Scala API,请参阅 DStream 和 PairDStreamFunctions。有关 Java API,请参阅 JavaDStream 和 JavaPairDStream。
DStream 上的输出操作
输出操作允许 DStream 的数据被推送到外部系统,如数据库或文件系统。由于输出操作允许转换后的数据被外部系统消费,它们会触发所有 DStream 转换的实际执行(类似于 RDD 的 Action)。目前,定义了以下输出操作:
| 输出操作 | 含义 |
|---|---|
| print() | 在运行流应用程序的驱动节点上打印 DStream 中每个数据批次的前十个元素。这对于开发和调试非常有用。 Python API 在 Python API 中,这被称为 pprint()。 |
| saveAsTextFiles(prefix, [suffix]) | 将此 DStream 的内容保存为文本文件。每个批次间隔的文件名根据 prefix 和 suffix 生成:"prefix-TIME_IN_MS[.suffix]"。 |
| saveAsObjectFiles(prefix, [suffix]) | 将此 DStream 的内容保存为序列化 Java 对象的 SequenceFiles。每个批次间隔的文件名根据 prefix 和 suffix 生成:"prefix-TIME_IN_MS[.suffix]"。Python API 这在 Python API 中不可用。 |
| saveAsHadoopFiles(prefix, [suffix]) | 将此 DStream 的内容保存为 Hadoop 文件。每个批次间隔的文件名根据 prefix 和 suffix 生成:"prefix-TIME_IN_MS[.suffix]"。 Python API 这在 Python API 中不可用。 |
| foreachRDD(func) | 最通用的输出运算符,它将函数 func 应用于从流生成的每个 RDD。此函数应将每个 RDD 中的数据推送到外部系统,例如将 RDD 保存到文件,或通过网络将其写入数据库。注意,函数 func 在运行流应用程序的驱动进程中执行,并且通常会在其中包含 RDD 操作,这将强制计算流式 RDD。 |
使用 foreachRDD 的设计模式
dstream.foreachRDD 是一个强大的原语,允许将数据发送到外部系统。但是,了解如何正确且高效地使用此原语非常重要。一些需要避免的常见错误如下。
通常,将数据写入外部系统需要创建一个连接对象(例如 TCP 连接到远程服务器)并使用它将数据发送到远程系统。为此,开发人员可能会不经意地尝试在 Spark 驱动程序上创建连接对象,然后尝试在 Spark worker 中使用它来保存 RDD 中的记录。例如(在 Scala 中):
def sendRecord(rdd):
connection = createNewConnection() # executed at the driver
rdd.foreach(lambda record: connection.send(record))
connection.close()
dstream.foreachRDD(sendRecord)dstream.foreachRDD { rdd =>
val connection = createNewConnection() // executed at the driver
rdd.foreach { record =>
connection.send(record) // executed at the worker
}
}dstream.foreachRDD(rdd -> {
Connection connection = createNewConnection(); // executed at the driver
rdd.foreach(record -> {
connection.send(record); // executed at the worker
});
});这是不正确的,因为这需要连接对象被序列化并从驱动程序发送到 worker。此类连接对象很难在机器之间传输。此错误可能表现为序列化错误(连接对象不可序列化)、初始化错误(连接对象需要在 worker 上初始化)等。正确的解决方案是在 worker 上创建连接对象。
但是,这可能导致另一个常见错误——为每条记录创建一个新连接。例如:
def sendRecord(record):
connection = createNewConnection()
connection.send(record)
connection.close()
dstream.foreachRDD(lambda rdd: rdd.foreach(sendRecord))dstream.foreachRDD { rdd =>
rdd.foreach { record =>
val connection = createNewConnection()
connection.send(record)
connection.close()
}
}dstream.foreachRDD(rdd -> {
rdd.foreach(record -> {
Connection connection = createNewConnection();
connection.send(record);
connection.close();
});
});通常,创建连接对象有时间和资源开销。因此,为每条记录创建和销毁连接对象可能会产生不必要的高开销,并会显著降低系统的整体吞吐量。一个更好的解决方案是使用 rdd.foreachPartition - 创建一个单一的连接对象,并使用该连接发送 RDD 分区中的所有记录。
def sendPartition(iter):
connection = createNewConnection()
for record in iter:
connection.send(record)
connection.close()
dstream.foreachRDD(lambda rdd: rdd.foreachPartition(sendPartition))dstream.foreachRDD { rdd =>
rdd.foreachPartition { partitionOfRecords =>
val connection = createNewConnection()
partitionOfRecords.foreach(record => connection.send(record))
connection.close()
}
}dstream.foreachRDD(rdd -> {
rdd.foreachPartition(partitionOfRecords -> {
Connection connection = createNewConnection();
while (partitionOfRecords.hasNext()) {
connection.send(partitionOfRecords.next());
}
connection.close();
});
});这将把连接创建开销摊销到许多记录上。
最后,这可以通过在多个 RDD/批次之间重用连接对象来进一步优化。可以维护一个连接对象池,当多个批次的 RDD 被推送到外部系统时可以重用这些连接对象,从而进一步减少开销。
def sendPartition(iter):
# ConnectionPool is a static, lazily initialized pool of connections
connection = ConnectionPool.getConnection()
for record in iter:
connection.send(record)
# return to the pool for future reuse
ConnectionPool.returnConnection(connection)
dstream.foreachRDD(lambda rdd: rdd.foreachPartition(sendPartition))dstream.foreachRDD { rdd =>
rdd.foreachPartition { partitionOfRecords =>
// ConnectionPool is a static, lazily initialized pool of connections
val connection = ConnectionPool.getConnection()
partitionOfRecords.foreach(record => connection.send(record))
ConnectionPool.returnConnection(connection) // return to the pool for future reuse
}
}dstream.foreachRDD(rdd -> {
rdd.foreachPartition(partitionOfRecords -> {
// ConnectionPool is a static, lazily initialized pool of connections
Connection connection = ConnectionPool.getConnection();
while (partitionOfRecords.hasNext()) {
connection.send(partitionOfRecords.next());
}
ConnectionPool.returnConnection(connection); // return to the pool for future reuse
});
});注意,池中的连接应该是按需惰性创建的,如果一段时间未使用则应超时。这实现了向外部系统发送数据的最高效方式。
其他需要记住的要点
-
DStreams 由输出操作惰性执行,就像 RDD 由 RDD 操作惰性执行一样。具体来说,DStream 输出操作内部的 RDD 操作强制处理接收到的数据。因此,如果你的应用程序没有任何输出操作,或者有像
dstream.foreachRDD()这样内部没有任何 RDD 操作的输出操作,那么将不会执行任何操作。系统只会简单地接收数据并将其丢弃。 -
默认情况下,输出操作一次执行一个。它们按照在应用程序中定义的顺序执行。
DataFrame 和 SQL 操作
你可以轻松地在流数据上使用 DataFrames 和 SQL 操作。你必须使用 StreamingContext 正在使用的 SparkContext 来创建 SparkSession。此外,必须以一种在驱动程序故障时可以重新启动的方式来执行此操作。这是通过创建一个惰性实例化的 SparkSession 单例实例来完成的。以下示例展示了这一点。它修改了前面的 单词计数示例,使用 DataFrames 和 SQL 生成单词计数。每个 RDD 都转换为 DataFrame,注册为临时表,然后使用 SQL 进行查询。
# Lazily instantiated global instance of SparkSession
def getSparkSessionInstance(sparkConf):
if ("sparkSessionSingletonInstance" not in globals()):
globals()["sparkSessionSingletonInstance"] = SparkSession \
.builder \
.config(conf=sparkConf) \
.getOrCreate()
return globals()["sparkSessionSingletonInstance"]
...
# DataFrame operations inside your streaming program
words = ... # DStream of strings
def process(time, rdd):
print("========= %s =========" % str(time))
try:
# Get the singleton instance of SparkSession
spark = getSparkSessionInstance(rdd.context.getConf())
# Convert RDD[String] to RDD[Row] to DataFrame
rowRdd = rdd.map(lambda w: Row(word=w))
wordsDataFrame = spark.createDataFrame(rowRdd)
# Creates a temporary view using the DataFrame
wordsDataFrame.createOrReplaceTempView("words")
# Do word count on table using SQL and print it
wordCountsDataFrame = spark.sql("select word, count(*) as total from words group by word")
wordCountsDataFrame.show()
except:
pass
words.foreachRDD(process)查看完整的 源代码。
/** DataFrame operations inside your streaming program */
val words: DStream[String] = ...
words.foreachRDD { rdd =>
// Get the singleton instance of SparkSession
val spark = SparkSession.builder.config(rdd.sparkContext.getConf).getOrCreate()
import spark.implicits._
// Convert RDD[String] to DataFrame
val wordsDataFrame = rdd.toDF("word")
// Create a temporary view
wordsDataFrame.createOrReplaceTempView("words")
// Do word count on DataFrame using SQL and print it
val wordCountsDataFrame =
spark.sql("select word, count(*) as total from words group by word")
wordCountsDataFrame.show()
}查看完整的 源代码。
/** Java Bean class for converting RDD to DataFrame */
public class JavaRow implements java.io.Serializable {
private String word;
public String getWord() {
return word;
}
public void setWord(String word) {
this.word = word;
}
}
...
/** DataFrame operations inside your streaming program */
JavaDStream<String> words = ...
words.foreachRDD((rdd, time) -> {
// Get the singleton instance of SparkSession
SparkSession spark = SparkSession.builder().config(rdd.sparkContext().getConf()).getOrCreate();
// Convert RDD[String] to RDD[case class] to DataFrame
JavaRDD<JavaRow> rowRDD = rdd.map(word -> {
JavaRow record = new JavaRow();
record.setWord(word);
return record;
});
DataFrame wordsDataFrame = spark.createDataFrame(rowRDD, JavaRow.class);
// Creates a temporary view using the DataFrame
wordsDataFrame.createOrReplaceTempView("words");
// Do word count on table using SQL and print it
DataFrame wordCountsDataFrame =
spark.sql("select word, count(*) as total from words group by word");
wordCountsDataFrame.show();
});查看完整的 源代码。
你还可以从不同的线程(即与正在运行的 StreamingContext 异步)对定义在流数据上的表运行 SQL 查询。只需确保你将 StreamingContext 设置为记住足够数量的流数据,以便查询可以运行。否则,不了解任何异步 SQL 查询的 StreamingContext 将在查询完成之前删除旧的流数据。例如,如果你想查询最后一个批次,但你的查询可能需要 5 分钟才能运行,那么请调用 streamingContext.remember(Minutes(5))(在 Scala 中,或在其他语言中等效)。
参阅 DataFrames 和 SQL 指南以了解有关 DataFrames 的更多信息。
MLlib 操作
你还可以轻松使用 MLlib 提供的机器学习算法。首先,有流式机器学习算法(例如 流式线性回归,流式 KMeans 等),它们既可以从流数据中学习,也可以将模型应用于流数据。除此之外,对于更大范围的机器学习算法,你可以离线学习学习模型(即使用历史数据),然后在线将模型应用于流数据。有关更多详细信息,请参阅 MLlib 指南。
缓存 / 持久化
与 RDD 类似,DStreams 也允许开发人员将流的数据持久化在内存中。也就是说,在 DStream 上使用 persist() 方法将自动将该 DStream 的每个 RDD 持久化在内存中。如果 DStream 中的数据将被多次计算(例如,对同一数据进行多次操作),这很有用。对于基于窗口的操作(如 reduceByWindow 和 reduceByKeyAndWindow)以及基于状态的操作(如 updateStateByKey),这隐式成立。因此,由基于窗口的操作生成的 DStreams 会自动持久化在内存中,无需开发人员调用 persist()。
对于通过网络接收数据的输入流(如 Kafka、套接字等),默认持久化级别设置为将数据复制到两个节点以实现容错。
注意,与 RDD 不同,DStreams 的默认持久化级别将数据序列化在内存中。这将在 性能调优 部分进一步讨论。有关不同持久化级别的更多信息,可以在 Spark 编程指南 中找到。
检查点 (Checkpointing)
流应用程序必须 24/7 运行,因此必须能够抵御与应用程序逻辑无关的故障(例如系统故障、JVM 崩溃等)。为此,Spark Streaming 需要将足够的信息 检查点 (checkpoint) 到容错存储系统中,以便从故障中恢复。有两种类型的数据被检查点。
- 元数据检查点 - 保存定义流计算的信息到 HDFS 等容错存储中。这用于从运行流应用程序驱动程序的节点的故障中恢复(稍后详细讨论)。元数据包括:
- 配置 - 用于创建流应用程序的配置。
- DStream 操作 - 定义流应用程序的一组 DStream 操作。
- 不完整批次 - 作业已排队但尚未完成的批次。
- 数据检查点 - 将生成的 RDD 保存到可靠存储中。这在某些跨多个批次组合数据的 有状态 转换中是必要的。在此类转换中,生成的 RDD 依赖于先前批次的 RDD,这会导致依赖链的长度随时间不断增加。为了避免恢复时间出现这种无限制的增加(与依赖链成正比),有状态转换的中间 RDD 会定期 检查点 到可靠存储(例如 HDFS)以切断依赖链。
总之,元数据检查点主要用于从驱动程序故障中恢复,而数据或 RDD 检查点即使在基础功能中也是必要的(如果使用了有状态转换)。
何时启用检查点
对于具有以下任何要求的应用程序,必须启用检查点:
- 使用有状态转换 - 如果应用程序中使用了
updateStateByKey或reduceByKeyAndWindow(带有反函数),则必须提供检查点目录以允许定期进行 RDD 检查点。 - 从运行应用程序的驱动程序故障中恢复 - 元数据检查点用于通过进度信息进行恢复。
注意,没有上述有状态转换的简单流应用程序可以在不启用检查点的情况下运行。在这种情况下,从驱动程序故障中恢复也将是部分的(可能会丢失一些已接收但未处理的数据)。这通常是可以接受的,许多人以这种方式运行 Spark Streaming 应用程序。对非 Hadoop 环境的支持预计在未来会得到改善。
如何配置检查点
可以通过在容错、可靠的文件系统(例如 HDFS、S3 等)中设置一个目录来启用检查点,检查点信息将保存到该目录中。这是通过使用 streamingContext.checkpoint(checkpointDirectory) 完成的。这将允许你使用上述有状态转换。此外,如果你想让应用程序从驱动程序故障中恢复,你应该重写你的流应用程序以具有以下行为:
- 当程序第一次启动时,它将创建一个新的 StreamingContext,设置所有流,然后调用 start()。
- 当程序在故障后重新启动时,它将从检查点目录中的检查点数据重新创建一个 StreamingContext。
通过使用 StreamingContext.getOrCreate,这种行为变得简单。用法如下:
# Function to create and setup a new StreamingContext
def functionToCreateContext():
sc = SparkContext(...) # new context
ssc = StreamingContext(...)
lines = ssc.socketTextStream(...) # create DStreams
...
ssc.checkpoint(checkpointDirectory) # set checkpoint directory
return ssc
# Get StreamingContext from checkpoint data or create a new one
context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext)
# Do additional setup on context that needs to be done,
# irrespective of whether it is being started or restarted
context. ...
# Start the context
context.start()
context.awaitTermination()如果 checkpointDirectory 存在,则上下文将从检查点数据重新创建。如果目录不存在(即第一次运行),则将调用函数 functionToCreateContext 来创建一个新上下文并设置 DStreams。查看 Python 示例 recoverable_network_wordcount.py。此示例将网络数据的单词计数附加到文件中。
你还可以通过使用 StreamingContext.getOrCreate(checkpointDirectory, None) 显式地从检查点数据创建 StreamingContext 并启动计算。
通过使用 StreamingContext.getOrCreate,这种行为变得简单。用法如下:
// Function to create and setup a new StreamingContext
def functionToCreateContext(): StreamingContext = {
val ssc = new StreamingContext(...) // new context
val lines = ssc.socketTextStream(...) // create DStreams
...
ssc.checkpoint(checkpointDirectory) // set checkpoint directory
ssc
}
// Get StreamingContext from checkpoint data or create a new one
val context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext _)
// Do additional setup on context that needs to be done,
// irrespective of whether it is being started or restarted
context. ...
// Start the context
context.start()
context.awaitTermination()如果 checkpointDirectory 存在,则上下文将从检查点数据重新创建。如果目录不存在(即第一次运行),则将调用函数 functionToCreateContext 来创建一个新上下文并设置 DStreams。查看 Scala 示例 RecoverableNetworkWordCount。此示例将网络数据的单词计数附加到文件中。
通过使用 JavaStreamingContext.getOrCreate,这种行为变得简单。用法如下:
// Create a factory object that can create and setup a new JavaStreamingContext
JavaStreamingContextFactory contextFactory = new JavaStreamingContextFactory() {
@Override public JavaStreamingContext create() {
JavaStreamingContext jssc = new JavaStreamingContext(...); // new context
JavaDStream<String> lines = jssc.socketTextStream(...); // create DStreams
...
jssc.checkpoint(checkpointDirectory); // set checkpoint directory
return jssc;
}
};
// Get JavaStreamingContext from checkpoint data or create a new one
JavaStreamingContext context = JavaStreamingContext.getOrCreate(checkpointDirectory, contextFactory);
// Do additional setup on context that needs to be done,
// irrespective of whether it is being started or restarted
context. ...
// Start the context
context.start();
context.awaitTermination();如果 checkpointDirectory 存在,则上下文将从检查点数据重新创建。如果目录不存在(即第一次运行),则将调用函数 contextFactory 来创建一个新上下文并设置 DStreams。查看 Java 示例 JavaRecoverableNetworkWordCount。此示例将网络数据的单词计数附加到文件中。
除了使用 getOrCreate 之外,还需要确保驱动程序进程在故障时自动重新启动。这只能由用于运行应用程序的部署基础结构来完成。这将在 部署 部分进一步讨论。
注意,RDD 的检查点会产生保存到可靠存储的成本。这可能会导致 RDD 被检查点的那些批次的处理时间增加。因此,需要仔细设置检查点间隔。在小批次大小(例如 1 秒)下,每批次检查点可能会显著降低操作吞吐量。相反,检查点太不频繁会导致血统和任务大小增加,这可能会产生不利影响。对于需要 RDD 检查点的有状态转换,默认间隔是批处理间隔的倍数,至少为 10 秒。它可以通过使用 dstream.checkpoint(checkpointInterval) 来设置。通常,5 到 10 个 DStream 滑动间隔的检查点间隔是一个很好的尝试设置。
累加器、广播变量和检查点
累加器 和 广播变量 无法从 Spark Streaming 中的检查点恢复。如果你启用了检查点并同时使用了 累加器 或 广播变量,则必须为 累加器 和 广播变量 创建惰性实例化的单例实例,以便它们可以在驱动程序在故障后重新启动时重新实例化。以下示例展示了这一点。
def getWordExcludeList(sparkContext):
if ("wordExcludeList" not in globals()):
globals()["wordExcludeList"] = sparkContext.broadcast(["a", "b", "c"])
return globals()["wordExcludeList"]
def getDroppedWordsCounter(sparkContext):
if ("droppedWordsCounter" not in globals()):
globals()["droppedWordsCounter"] = sparkContext.accumulator(0)
return globals()["droppedWordsCounter"]
def echo(time, rdd):
# Get or register the excludeList Broadcast
excludeList = getWordExcludeList(rdd.context)
# Get or register the droppedWordsCounter Accumulator
droppedWordsCounter = getDroppedWordsCounter(rdd.context)
# Use excludeList to drop words and use droppedWordsCounter to count them
def filterFunc(wordCount):
if wordCount[0] in excludeList.value:
droppedWordsCounter.add(wordCount[1])
False
else:
True
counts = "Counts at time %s %s" % (time, rdd.filter(filterFunc).collect())
wordCounts.foreachRDD(echo)查看完整的 源代码。
object WordExcludeList {
@volatile private var instance: Broadcast[Seq[String]] = null
def getInstance(sc: SparkContext): Broadcast[Seq[String]] = {
if (instance == null) {
synchronized {
if (instance == null) {
val wordExcludeList = Seq("a", "b", "c")
instance = sc.broadcast(wordExcludeList)
}
}
}
instance
}
}
object DroppedWordsCounter {
@volatile private var instance: LongAccumulator = null
def getInstance(sc: SparkContext): LongAccumulator = {
if (instance == null) {
synchronized {
if (instance == null) {
instance = sc.longAccumulator("DroppedWordsCounter")
}
}
}
instance
}
}
wordCounts.foreachRDD { (rdd: RDD[(String, Int)], time: Time) =>
// Get or register the excludeList Broadcast
val excludeList = WordExcludeList.getInstance(rdd.sparkContext)
// Get or register the droppedWordsCounter Accumulator
val droppedWordsCounter = DroppedWordsCounter.getInstance(rdd.sparkContext)
// Use excludeList to drop words and use droppedWordsCounter to count them
val counts = rdd.filter { case (word, count) =>
if (excludeList.value.contains(word)) {
droppedWordsCounter.add(count)
false
} else {
true
}
}.collect().mkString("[", ", ", "]")
val output = "Counts at time " + time + " " + counts
})查看完整的 源代码。
class JavaWordExcludeList {
private static volatile Broadcast<List<String>> instance = null;
public static Broadcast<List<String>> getInstance(JavaSparkContext jsc) {
if (instance == null) {
synchronized (JavaWordExcludeList.class) {
if (instance == null) {
List<String> wordExcludeList = Arrays.asList("a", "b", "c");
instance = jsc.broadcast(wordExcludeList);
}
}
}
return instance;
}
}
class JavaDroppedWordsCounter {
private static volatile LongAccumulator instance = null;
public static LongAccumulator getInstance(JavaSparkContext jsc) {
if (instance == null) {
synchronized (JavaDroppedWordsCounter.class) {
if (instance == null) {
instance = jsc.sc().longAccumulator("DroppedWordsCounter");
}
}
}
return instance;
}
}
wordCounts.foreachRDD((rdd, time) -> {
// Get or register the excludeList Broadcast
Broadcast<List<String>> excludeList = JavaWordExcludeList.getInstance(new JavaSparkContext(rdd.context()));
// Get or register the droppedWordsCounter Accumulator
LongAccumulator droppedWordsCounter = JavaDroppedWordsCounter.getInstance(new JavaSparkContext(rdd.context()));
// Use excludeList to drop words and use droppedWordsCounter to count them
String counts = rdd.filter(wordCount -> {
if (excludeList.value().contains(wordCount._1())) {
droppedWordsCounter.add(wordCount._2());
return false;
} else {
return true;
}
}).collect().toString();
String output = "Counts at time " + time + " " + counts;
}查看完整的 源代码。
部署应用程序
本节讨论部署 Spark Streaming 应用程序的步骤。
要求
要运行 Spark Streaming 应用程序,你需要具备以下条件:
-
具有集群管理器的集群 - 这是任何 Spark 应用程序的一般要求,并在 部署指南 中进行了详细讨论。
-
打包应用程序 JAR - 你必须将流应用程序编译为 JAR。如果你使用
spark-submit启动应用程序,则无需在 JAR 中提供 Spark 和 Spark Streaming。但是,如果你的应用程序使用了 高级源(例如 Kafka),则必须将它们链接到的额外构件及其依赖项打包在用于部署应用程序的 JAR 中。例如,使用KafkaUtils的应用程序必须在其应用程序 JAR 中包含spark-streaming-kafka-0-10_2.13及其所有传递依赖项。 -
为执行器配置足够的内存 - 由于接收到的数据必须存储在内存中,因此必须将执行器配置为具有足够的内存来保存接收到的数据。注意,如果你正在执行 10 分钟的窗口操作,系统必须在内存中至少保留过去 10 分钟的数据。因此,应用程序的内存需求取决于其中使用的操作。
-
配置检查点 - 如果流应用程序需要,则必须将 Hadoop API 兼容的容错存储(例如 HDFS、S3 等)中的目录配置为检查点目录,并以可以使用检查点信息进行故障恢复的方式编写流应用程序。有关详细信息,请参阅 检查点 部分。
- 配置应用程序驱动程序的自动重启 - 为了从驱动程序故障中自动恢复,用于运行流应用程序的部署基础结构必须监控驱动程序进程,并在驱动程序故障时重新启动它。不同的 集群管理器 有不同的工具来实现这一点。
- Spark Standalone - 可以提交 Spark 应用程序驱动程序以在 Spark Standalone 集群中运行(参阅 集群部署模式),也就是说,应用程序驱动程序本身运行在其中一个工作节点上。此外,可以指示 Standalone 集群管理器 监督 驱动程序,并在驱动程序由于非零退出代码或运行驱动程序的节点故障而失败时重新启动它。有关详细信息,请参阅 Spark Standalone 指南 中的 集群模式 和 监督。
- YARN - Yarn 支持类似的自动重启应用程序的机制。请参阅 YARN 文档以了解更多详细信息。
-
配置预写日志 (Write-ahead Logs) - 自 Spark 1.2 起,我们引入了 预写日志 以实现强大的容错保证。如果启用,从接收器接收到的所有数据都会写入配置检查点目录中的预写日志中。这可以防止驱动程序恢复时的数据丢失,从而确保零数据丢失(在 容错语义 部分详细讨论)。可以通过将 配置参数
spark.streaming.receiver.writeAheadLog.enable设置为true来启用此功能。但是,这些更强的语义可能会以牺牲单个接收器的接收吞吐量为代价。这可以通过运行 更多并行接收器 来增加总吞吐量来纠正。此外,建议在启用预写日志时禁用 Spark 内接收数据的复制,因为日志已存储在复制存储系统中。这可以通过将输入流的存储级别设置为StorageLevel.MEMORY_AND_DISK_SER来完成。在使用 S3(或任何不支持刷新操作的文件系统)进行 预写日志 时,请记住启用spark.streaming.driver.writeAheadLog.closeFileAfterWrite和spark.streaming.receiver.writeAheadLog.closeFileAfterWrite。有关更多详细信息,请参阅 Spark Streaming 配置。注意,当启用 I/O 加密时,Spark 不会加密写入预写日志的数据。如果需要预写日志数据的加密,它应存储在原生支持加密的文件系统中。 - 设置最大接收速率 - 如果集群资源不足以让流应用程序以接收数据的速度处理数据,则可以通过设置记录/秒的最大速率限制来限制接收器。请参阅接收器的 配置参数
spark.streaming.receiver.maxRate和 Direct Kafka 方法的spark.streaming.kafka.maxRatePerPartition。在 Spark 1.5 中,我们引入了一种称为 背压 (backpressure) 的功能,它消除了设置此速率限制的必要,因为 Spark Streaming 会自动计算速率限制并在处理条件发生变化时动态调整它们。可以通过将 配置参数spark.streaming.backpressure.enabled设置为true来启用此背压功能。
升级应用程序代码
如果需要使用新的应用程序代码升级正在运行的 Spark Streaming 应用程序,则有两种可能的机制。
-
升级后的 Spark Streaming 应用程序与现有应用程序并行启动和运行。一旦新的应用程序(接收与旧应用程序相同的数据)预热并准备就绪,旧的应用程序就可以关闭。注意,这对于支持将数据发送到两个目的地(即之前和升级后的应用程序)的数据源来说是可以做到的。
-
现有应用程序被优雅地关闭(有关优雅关闭选项,请参阅
StreamingContext.stop(...)或JavaStreamingContext.stop(...)),这确保了在关闭之前已接收的数据得到完全处理。然后可以启动升级后的应用程序,它将从前一个应用程序停止的地方开始处理。注意,这只能在支持源端缓冲(如 Kafka)的输入源中完成,因为在之前应用程序关闭而升级后的应用程序尚未启动时,数据需要被缓冲。并且无法从预升级代码的先前检查点信息重新启动。检查点信息本质上包含序列化的 Scala/Java/Python 对象,尝试使用新的、修改后的类反序列化对象可能会导致错误。在这种情况下,要么使用不同的检查点目录启动升级后的应用程序,要么删除先前的检查点目录。
监控应用程序
除了 Spark 的 监控功能 之外,还有特定于 Spark Streaming 的附加功能。当使用 StreamingContext 时,Spark Web UI 会显示一个额外的 Streaming 选项卡,其中显示有关运行中接收器(接收器是否处于活动状态、接收到的记录数、接收器错误等)和已完成批次(批处理时间、排队延迟等)的统计信息。这可用于监控流应用程序的进度。
Web UI 中的以下两个指标尤为重要:
- 处理时间 (Processing Time) - 处理每个数据批次所需的时间。
- 调度延迟 (Scheduling Delay) - 批次在队列中等待之前批次处理完成的时间。
如果批处理时间持续超过批处理间隔和/或排队延迟不断增加,则表明系统无法以生成数据的速度处理批次,并且正在落后。在这种情况下,请考虑 减少 批处理时间。
Spark Streaming 程序的进度也可以使用 StreamingListener 接口进行监控,该接口允许你获取接收器状态和处理时间。注意,这是一个开发人员 API,预计在未来会得到改进(即报告更多信息)。
性能调优
在集群上获得 Spark Streaming 应用程序的最佳性能需要进行一些调整。本节解释了可以调整以提高应用程序性能的多个参数和配置。在宏观层面上,你需要考虑两件事:
-
通过有效地使用集群资源来减少每个数据批次的处理时间。
-
设置正确的批处理大小,以便数据批次可以以接收它们的速度进行处理(即数据处理跟上数据摄入)。
减少批处理时间
Spark 中可以进行许多优化以最大限度地减少每个批次的处理时间。这些已在 调优指南 中进行了详细讨论。本节重点介绍了其中一些最重要的优化。
数据接收的并行度
通过网络接收数据(如 Kafka、套接字等)要求数据被反序列化并存储在 Spark 中。如果数据接收成为系统中的瓶颈,则考虑并行化数据接收。注意,每个输入 DStream 创建一个单一的接收器(运行在 worker 机器上),该接收器接收单个数据流。因此,接收多个数据流可以通过创建多个输入 DStreams 并配置它们从源接收数据流的不同分区来实现。例如,接收两个主题数据的单个 Kafka 输入 DStream 可以拆分为两个 Kafka 输入流,每个流仅接收一个主题。这将运行两个接收器,从而允许并行接收数据,进而增加总吞吐量。可以将这些多个 DStreams 连接在一起以创建一个单一的 DStream。然后,可以对统一后的流应用之前应用于单个输入 DStream 的转换。这是按照以下方式完成的:
numStreams = 5
kafkaStreams = [KafkaUtils.createStream(...) for _ in range (numStreams)]
unifiedStream = streamingContext.union(*kafkaStreams)
unifiedStream.pprint()val numStreams = 5
val kafkaStreams = (1 to numStreams).map { i => KafkaUtils.createStream(...) }
val unifiedStream = streamingContext.union(kafkaStreams)
unifiedStream.print()int numStreams = 5;
List<JavaPairDStream<String, String>> kafkaStreams = new ArrayList<>(numStreams);
for (int i = 0; i < numStreams; i++) {
kafkaStreams.add(KafkaUtils.createStream(...));
}
JavaPairDStream<String, String> unifiedStream = streamingContext.union(kafkaStreams.get(0), kafkaStreams.subList(1, kafkaStreams.size()));
unifiedStream.print();另一个应该考虑的参数是接收器的块间隔 (block interval),该间隔由 配置参数 spark.streaming.blockInterval 决定。对于大多数接收器,接收到的数据在存储到 Spark 内存之前会被合并成数据块。每批次中的块数决定了在 map 类转换中用于处理接收到的数据的任务数。每批次每个接收器的任务数将大约为(批处理间隔 / 块间隔)。例如,200 毫秒的块间隔将为每 2 秒的批次创建 10 个任务。如果任务数太少(即少于每台机器的核心数),那么它将是低效的,因为所有可用核心都不会被用于处理数据。要增加给定批处理间隔的任务数,请减小块间隔。但是,块间隔的建议最小值为约 50 毫秒,低于此值可能会出现任务启动开销问题。
接收具有多个输入流/接收器的数据的替代方法是显式重分区输入数据流(使用 inputStream.repartition(<number of partitions>))。这会在进一步处理之前将接收到的数据批次分发到集群中指定数量的机器上。
对于直接流,请参考 Spark Streaming + Kafka 集成指南
数据处理的并行度
如果计算的任何阶段所使用的并行任务数量不够高,集群资源可能会被利用不足。例如,对于像 reduceByKey 和 reduceByKeyAndWindow 这样的分布式 reduce 操作,默认的并行任务数量由 spark.default.parallelism 配置属性控制。你可以将并行度作为参数传递(参见 PairDStreamFunctions 文档),或者设置 spark.default.parallelism 配置属性来更改默认值。
数据序列化
通过调整序列化格式可以降低数据序列化的开销。在流处理中,有两种类型的数据需要序列化。
-
输入数据:默认情况下,通过接收器(Receivers)接收的输入数据以 StorageLevel.MEMORY_AND_DISK_SER_2 存储在执行器(executors)的内存中。这意味着,为了减少 GC 开销,数据会被序列化为字节,并进行复制以容忍执行器故障。此外,数据首先保留在内存中,只有在内存不足以容纳流计算所需的所有输入数据时,才会溢出到磁盘。这种序列化显然存在开销——接收器必须反序列化接收到的数据,并使用 Spark 的序列化格式重新序列化它。
-
流式操作生成的持久化 RDD:流式计算生成的 RDD 可能会持久化在内存中。例如,窗口操作会将数据持久化在内存中,因为它们会被多次处理。然而,与 Spark Core 默认的 StorageLevel.MEMORY_ONLY 不同,流式计算生成的持久化 RDD 默认使用 StorageLevel.MEMORY_ONLY_SER(即序列化)进行持久化,以最小化 GC 开销。
在这两种情况下,使用 Kryo 序列化都可以减少 CPU 和内存开销。有关更多详细信息,请参阅 Spark 调优指南。对于 Kryo,请考虑注册自定义类,并禁用对象引用跟踪(参见配置指南中与 Kryo 相关的配置)。
在流式应用所需保留的数据量不大的特定情况下,将数据(上述两种类型)作为反序列化对象进行持久化,而不产生过多的 GC 开销可能是可行的。例如,如果你使用的批处理间隔为几秒钟且没有窗口操作,那么你可以尝试通过相应地显式设置存储级别来禁用持久化数据中的序列化。这将减少由于序列化导致的 CPU 开销,可能在不产生过多 GC 开销的情况下提高性能。
设置合理的批处理间隔
为了使运行在集群上的 Spark Streaming 应用保持稳定,系统应该能够以与接收数据相同的速度处理数据。换句话说,批次数据的处理速度应与生成速度一致。可以通过在流处理 Web UI 中监控处理时间来判断应用是否满足这一点,其中批处理时间应小于批处理间隔。
根据流计算的性质,所使用的批处理间隔可能会对应用在固定集群资源下所能维持的数据速率产生重大影响。例如,让我们考虑早期的 WordCountNetwork 示例。对于特定的数据速率,系统可能能够跟上每 2 秒报告一次词频(即 2 秒的批处理间隔),但无法跟上每 500 毫秒一次。因此,批处理间隔的设置需要确保生产环境中的预期数据速率能够被维持。
确定应用正确批处理大小的一个好方法是使用保守的批处理间隔(例如 5-10 秒)和低数据速率进行测试。为了验证系统是否能够跟上数据速率,你可以检查每个已处理批次的端到端延迟值(查看 Spark 驱动程序的 log4j 日志中的“Total delay”,或使用 StreamingListener 接口)。如果延迟保持在与批处理大小相当的水平,则系统是稳定的。否则,如果延迟持续增加,则意味着系统无法跟上,因此是不稳定的。一旦你对稳定配置有了概念,就可以尝试提高数据速率和/或减小批处理大小。注意,由于临时数据速率增加导致的短暂延迟增加是可以接受的,只要延迟回落到较低值(即小于批处理大小)即可。
内存调优
关于调整 Spark 应用的内存使用和 GC 行为,调优指南中已有详细讨论。强烈建议阅读该指南。在本节中,我们将专门讨论 Spark Streaming 应用背景下的一些调优参数。
Spark Streaming 应用所需的集群内存量在很大程度上取决于所使用的转换类型。例如,如果你想对过去 10 分钟的数据使用窗口操作,那么你的集群应有足够的内存来容纳 10 分钟的数据。或者,如果你想对大量键使用 updateStateByKey,那么所需的内存会很高。相反,如果你只是进行简单的 map-filter-store 操作,那么所需的内存则较低。
通常情况下,由于通过接收器接收的数据以 StorageLevel.MEMORY_AND_DISK_SER_2 存储,无法放入内存的数据会溢出到磁盘。这可能会降低流应用的处理性能,因此建议根据流应用的需求提供充足的内存。最好在小规模环境下观察内存使用情况,并进行相应估算。
内存调优的另一个方面是垃圾回收(GC)。对于需要低延迟的流应用,由 JVM 垃圾回收引起的大停顿是不可取的。
有几个参数可以帮助你调整内存使用和 GC 开销。
-
DStream 的持久化级别:如前文在 数据序列化 部分所述,输入数据和 RDD 默认以序列化字节形式持久化。与反序列化持久化相比,这减少了内存使用和 GC 开销。启用 Kryo 序列化进一步减小了序列化大小和内存使用。通过压缩(参见 Spark 配置
spark.rdd.compress)可以以 CPU 时间为代价进一步减少内存使用。 -
清理旧数据:默认情况下,DStream 转换生成的所有输入数据和持久化 RDD 都会被自动清理。Spark Streaming 根据所使用的转换决定何时清理数据。例如,如果你使用 10 分钟的窗口操作,那么 Spark Streaming 将保留最近 10 分钟的数据,并主动丢弃更早的数据。通过设置
streamingContext.remember,可以将数据保留更长的时间(例如为了交互式查询旧数据)。 -
其他技巧:为了进一步降低 GC 开销,这里还有一些建议可供尝试。
- 使用
OFF_HEAP存储级别持久化 RDD。详细信息请参阅 Spark 编程指南。 - 使用更多具有较小堆内存的执行器(executors)。这将减少每个 JVM 堆内的 GC 压力。
- 使用
需要记住的要点
-
一个 DStream 与单个接收器关联。为了实现读取并行性,需要创建多个接收器,即多个 DStream。接收器在执行器内运行,占用一个核心。在预定接收器插槽后,请确保有足够的内核用于处理,即
spark.cores.max应将接收器插槽考虑在内。接收器以轮询方式分配给执行器。 -
当数据从流源接收时,接收器会创建数据块。每隔 blockInterval 毫秒会生成一个新的数据块。在 batchInterval 期间创建 N 个数据块,其中 N = batchInterval/blockInterval。这些块由当前执行器的 BlockManager 分发到其他执行器的 BlockManager。之后,运行在驱动程序上的 Network Input Tracker 会获知块的位置以进行进一步处理。
-
在 batchInterval 期间创建的块会在驱动程序上创建一个 RDD。batchInterval 期间生成的块是 RDD 的分区。每个分区都是 Spark 中的一个任务。blockInterval == batchInterval 意味着只创建一个分区,并且通常会在本地处理。
-
无论 blockInterval 如何,块上的 map 任务都会在拥有这些块的执行器(接收块的执行器和复制块的执行器)上处理,除非触发了非本地调度。更大的 blockInterval 意味着更大的块。较高的
spark.locality.wait值增加了在本地节点处理块的机会。需要在这两个参数之间找到平衡,以确保更大的块能够被本地处理。 -
你可以通过调用
inputDstream.repartition(n)来定义分区数量,而不是依赖 batchInterval 和 blockInterval。这会随机重排 RDD 中的数据以创建 n 个分区。是的,为了获得更高的并行度。虽然这会以 shuffle 为代价。RDD 的处理由驱动程序的任务调度器作为作业进行调度。在给定时间点,只有一个作业处于活动状态。因此,如果一个作业正在执行,其他作业将被排队。 -
如果你有两个 DStream,将会形成两个 RDD,并创建两个作业,它们将依次调度。为了避免这种情况,你可以合并两个 DStream。这将确保为两个 DStream 的 RDD 形成一个单一的 unionRDD。此 unionRDD 然后被视为单个作业。但是,RDD 的分区不受此影响。
-
如果批处理时间超过了 batchInterval,那么显然接收器的内存将开始填满,最终会抛出异常(很可能是 BlockNotFoundException)。目前无法暂停接收器。可以使用 SparkConf 配置
spark.streaming.receiver.maxRate来限制接收器的速率。
容错语义
在本节中,我们将讨论 Spark Streaming 应用在发生故障时的行为。
背景
要理解 Spark Streaming 提供的语义,让我们回顾一下 Spark RDD 的基本容错语义。
- RDD 是一个不可变的、确定性可重计算的分布式数据集。每个 RDD 都记得用于在容错输入数据集上创建它的一系列确定性操作。
- 如果 RDD 的任何分区由于工作节点故障而丢失,则可以使用操作的血缘关系从原始容错数据集重新计算该分区。
- 假设所有的 RDD 转换都是确定的,无论 Spark 集群中发生什么故障,最终转换后的 RDD 中的数据将始终相同。
Spark 在 HDFS 或 S3 等容错文件系统上的数据上运行。因此,从容错数据生成的所有 RDD 也是容错的。然而,Spark Streaming 的情况并非如此,因为在大多数情况下,数据是通过网络接收的(除非使用 fileStream)。为了使所有生成的 RDD 达到相同的容错属性,接收到的数据会在集群中工作节点上的多个 Spark 执行器之间进行复制(默认复制因子为 2)。这导致系统中有两种类型的数据需要在发生故障时进行恢复:
- 接收并复制的数据 - 此数据可以抵御单个工作节点的故障,因为它的一份副本存在于其他节点上。
- 已接收但缓冲用于复制的数据 - 由于此数据未被复制,恢复此数据的唯一方法是从源头重新获取它。
此外,我们应该关注两种类型的故障:
- 工作节点故障 - 运行执行器的任何工作节点都可能发生故障,该节点上的所有内存数据都将丢失。如果故障节点上运行了任何接收器,则其缓冲的数据将会丢失。
- 驱动节点故障 - 如果运行 Spark Streaming 应用的驱动节点发生故障,显然 SparkContext 会丢失,所有带有内存数据的执行器也会丢失。
有了这些基本知识,让我们来了解 Spark Streaming 的容错语义。
定义
流系统的语义通常通过每条记录被系统处理的次数来捕获。系统在所有可能的运行条件下(尽管有故障等)可以提供三种类型的保证:
- 最多一次(At most once):每条记录要么被处理一次,要么根本不处理。
- 至少一次(At least once):每条记录将被处理一次或多次。这比最多一次更强,因为它确保数据不会丢失。但可能会出现重复。
- 恰好一次(Exactly once):每条记录将被恰好处理一次 - 数据不会丢失,也不会被处理多次。这显然是三种保证中最强的一种。
基本语义
在任何流处理系统中,概括地说,处理数据有三个步骤。
-
接收数据:数据通过接收器或其他方式从源头接收。
-
转换数据:接收到的数据使用 DStream 和 RDD 转换进行转换。
-
推送数据:最终转换后的数据被推送到文件系统、数据库、仪表板等外部系统。
如果流应用要实现端到端的恰好一次保证,那么每个步骤都必须提供恰好一次保证。也就是说,每条记录必须被恰好接收一次、恰好转换一次,并被恰好推送至下游系统一次。让我们在 Spark Streaming 的背景下了解这些步骤的语义。
-
接收数据:不同的输入源提供不同的保证。这将在下一小节中详细讨论。
-
转换数据:得益于 RDD 提供的保证,所有已接收的数据都将被恰好一次处理。即使发生故障,只要接收到的输入数据可访问,最终转换后的 RDD 内容将始终相同。
-
推送数据:输出操作默认确保至少一次语义,因为它取决于输出操作的类型(是否幂等)以及下游系统的语义(是否支持事务)。但用户可以实现自己的事务机制来达到恰好一次语义。本节稍后将详细讨论。
已接收数据的语义
不同的输入源提供不同的保证,范围从至少一次到恰好一次。阅读以了解更多详情。
对于文件
如果所有输入数据都已经存在于 HDFS 等容错文件系统中,Spark Streaming 总能从任何故障中恢复并处理所有数据。这提供了恰好一次语义,意味着无论发生什么故障,所有数据都将被恰好处理一次。
对于基于接收器的源
对于基于接收器的输入源,容错语义取决于故障场景和接收器的类型。正如我们前面讨论的,有两种类型的接收器
- 可靠接收器(Reliable Receiver) - 这些接收器仅在确保接收到的数据已被复制后才确认可靠源。如果此类接收器发生故障,源将不会收到针对缓冲(未复制)数据的确认。因此,如果接收器重启,源将重新发送数据,且不会因故障而丢失任何数据。
- 不可靠接收器(Unreliable Receiver) - 此类接收器不发送确认,因此在因工作节点或驱动故障而失败时可能丢失数据。
根据所使用的接收器类型,我们实现以下语义。如果工作节点发生故障,可靠接收器不会丢失数据。使用不可靠接收器时,已接收但未复制的数据可能会丢失。如果驱动节点发生故障,除了上述损失外,过去所有在内存中接收并复制的数据都将丢失。这将影响有状态转换的结果。
为了避免过去接收到的数据丢失,Spark 1.2 引入了预写日志(Write Ahead Logs),它将接收到的数据保存到容错存储中。在启用预写日志且使用可靠接收器的情况下,数据不会丢失。就语义而言,它提供了至少一次保证。
下表总结了故障情况下的语义
| 部署场景 | 工作节点故障 | 驱动节点故障 |
|---|---|---|
|
Spark 1.1 或更早版本, 或 Spark 1.2 或更高版本且未启用预写日志 |
不可靠接收器丢失缓冲数据 可靠接收器零数据丢失 至少一次语义 |
不可靠接收器丢失缓冲数据 所有接收器的过去数据丢失 语义未定义 |
| Spark 1.2 或更高版本且启用预写日志 | 可靠接收器零数据丢失 至少一次语义 |
可靠接收器和文件零数据丢失 至少一次语义 |
使用 Kafka Direct API
在 Spark 1.3 中,我们引入了新的 Kafka Direct API,它可以确保所有 Kafka 数据都能被 Spark Streaming 恰好一次地接收。在此基础上,如果你实现了恰好一次的输出操作,就可以实现端到端的恰好一次保证。这种方法在 Kafka 集成指南中作了进一步讨论。
输出操作的语义
输出操作(如 foreachRDD)具有至少一次语义,也就是说,在工作节点发生故障的情况下,转换后的数据可能会多次写入外部实体。虽然这对于使用 saveAs***Files 操作保存到文件系统是可以接受的(因为文件只会以相同的数据被覆盖),但可能需要额外的努力才能实现恰好一次语义。有两种方法。
-
幂等更新:多次尝试总是写入相同的数据。例如,
saveAs***Files总是将相同的数据写入生成的文件。 -
事务性更新:所有更新都以事务方式进行,确保更新原子地发生且仅发生一次。实现此目的的一种方法如下:
- 使用批处理时间(在
foreachRDD中可用)和 RDD 的分区索引创建一个标识符。此标识符唯一地标识流应用中的一个数据块(blob)。 -
使用该标识符将数据块事务性地更新到外部系统(即恰好一次,原子地)。也就是说,如果标识符尚未提交,则原子地提交分区数据和标识符。否则,如果已经提交,则跳过更新。
dstream.foreachRDD { (rdd, time) => rdd.foreachPartition { partitionIterator => val partitionId = TaskContext.get.partitionId() val uniqueId = generateUniqueId(time.milliseconds, partitionId) // use this uniqueId to transactionally commit the data in partitionIterator } }
- 使用批处理时间(在