Article

数据计算 Spark Streaming

更新于:2026-07-13

第一章:Spark Streaming 概述

1.1 什么是流处理

概念名称说明注意事项
流处理(Stream Processing)一种对连续不断生成的数据流进行实时或近实时处理的计算范式。与批处理不同,流处理系统持续接收、处理并输出结果,适用于日志分析、监控、推荐等场景。流处理强调低延迟,但不一定要求”毫秒级”响应;需权衡吞吐量与延迟。
数据流(Data Stream)无限序列的数据记录,按时间顺序持续到达,例如用户点击流、传感器数据、交易日志等。数据流通常是无界的(unbounded),不能一次性加载到内存中。
实时处理(Real-time Processing)在数据产生后极短时间内完成处理并返回结果的过程。Spark Streaming 属于”近实时”(near real-time),延迟通常在秒级。严格意义上的”实时”系统(hard real-time)要求确定性延迟,Spark Streaming 不属于此类。
批处理 vs 流处理批处理处理有界数据集(如 HDFS 上的文件),而流处理处理无界数据流。传统批处理延迟高,流处理延迟低。Spark Streaming 通过微批(micro-batch)方式统一了批与流的处理模型。

1.2 Spark Streaming 简介与核心思想

概念名称说明注意事项
Spark StreamingSpark 的流式计算模块,用于构建可扩展、容错的实时数据处理应用。它将流数据划分为一系列小的批处理任务,在 Spark 引擎上执行。基于 Spark Core 构建,可无缝集成 Spark SQL、MLlib、GraphX。
微批处理(Micro-batch Processing)核心思想是将连续的数据流切分为时间间隔较短的”批次”(如 1 秒),每个批次作为一个 RDD 进行处理。所有操作在批次级别执行,最小延迟为批处理间隔。
批处理间隔(Batch Interval)用户定义的时间间隔,决定数据被切分的频率(如 500ms、1s、5s)。影响延迟与吞吐量。设置过短会导致调度开销大;过长则延迟高。建议从 1s 起调优。
事件时间 vs 处理时间处理时间是数据被系统处理的时间;事件时间是数据实际发生的时间。Spark Streaming 默认使用处理时间。若需事件时间语义,推荐升级到 Structured Streaming。

1.3 DStream(离散化流)模型

概念名称说明注意事项
DStream(Discretized Stream)Spark Streaming 的基本抽象,表示一个连续的数据流。它是由时间序列组织的 RDD 序列,每个 RDD 包含一个时间间隔内的数据。DStream 是对 RDD 的封装,所有操作最终转化为 RDD 操作。
DStream 的生成方式可通过外部输入源(如 Kafka、Socket)或现有 DStream 经转换得到。接收器(Receiver)运行在 Executor 中,负责数据拉取。
DStream 的依赖关系每个 DStream 对象维护对父 DStream 的依赖,形成有向无环图(DAG),用于容错和调度。类似于 RDD 的 lineage 机制,支持故障恢复。
DStream 的生命周期从数据源开始,经过一系列转换(Transformation),最后通过输出操作(Output Operation)触发执行。DStream 本身是惰性的,只有遇到输出操作才会启动计算。

1.4 Spark Streaming 与其他流处理框架对比(如 Flink、Storm)

框架说明注意事项
Apache Storm真正的流处理系统,采用”记录级”处理,延迟可低至毫秒级。编程模型较复杂,状态管理需自行实现,运维成本高。
Apache Flink原生流处理引擎,支持事件时间、窗口、状态管理,提供精确一次(exactly-once)语义。学习曲线较陡,资源调度不如 Spark 生态集成紧密。
Spark Streaming基于微批处理,延迟在秒级,易于与 Spark 生态集成,开发成本低。不适合超低延迟场景;事件时间支持有限(DStream API)。
Kafka Streams轻量级库,直接嵌入 Java/Scala 应用,适合 Kafka 数据源的流处理。功能相对有限,不适合大规模复杂计算。
选择建议• 要求低延迟(<100ms):选 Flink 或 Storm
• 已有 Spark 生态:选 Spark Streaming
• 简单 Kafka 处理:选 Kafka Streams
Spark Streaming 正逐步被 Structured Streaming 取代,新项目建议优先考虑 Structured Streaming。

第二章:开发环境搭建与入门程序

2.1 环境准备(Scala、Spark、依赖配置)

组件说明注意事项
Scala 版本Spark 通常绑定特定 Scala 版本(如 Spark 3.x 支持 Scala 2.12)。需确保本地 Scala 版本匹配。推荐使用 Scala 2.12,避免版本冲突。
Spark 安装可从官网下载预编译版本,或通过包管理工具(如 Homebrew、apt)安装。需配置 SPARK_HOMEPATH本地模式无需集群,但生产环境需部署 Spark 集群(Standalone/YARN/K8s)。
Java 环境Spark 基于 JVM,需安装 JDK 8 或 11(推荐 OpenJDK)。确保 JAVA_HOME 正确设置。
依赖管理使用 Maven 或 SBT 管理项目依赖,需引入 spark-streaming_2.12 及其他连接器(如 Kafka)。注意 Spark 和依赖库的版本兼容性。

2.2 Maven/Gradle 项目构建

Maven 示例配置(pom.xml 片段)

配置项说明示例值
groupId项目组织标识com.example.streaming
artifactId项目名称spark-streaming-demo
scala.versionScala 版本2.12.15
spark.versionSpark 版本3.5.0
依赖项:Spark Streaming核心流处理库<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming_2.12</artifactId>
<version>3.5.0</version>
</dependency>
依赖项:Spark Core(自动引入)基础执行引擎通常无需显式声明,由 spark-streaming 传递依赖引入。
打包插件推荐使用 maven-shade-plugin 打包成 fat jar避免运行时类路径缺失问题。

注意事项:

  • 确保 <artifactId> 中的 Scala 版本(如 _2.12)与本地环境一致。
  • 使用 mvn compile 编译,mvn package 打包。

2.3 第一个 Spark Streaming 程序:WordCount(基于 Socket)

程序逻辑说明

概念说明注意事项
数据源使用 socketTextStream 从本地 TCP 端口(如 9999)读取文本流。需提前启动 Netcat 或 Python 脚本发送数据。
逻辑流程1. 创建 StreamingContext
2. 定义 DStream 从 socket 读取
3. 分词、映射为 (word, 1)
4. 聚合统计 word count
5. 输出结果
每个批次输出一次结果。
启动方式调用 ssc.start() 开始接收数据,ssc.awaitTermination() 等待终止信号。程序不会自动退出,需手动中断(Ctrl+C)。

核心方法与用途(Scala 示例)

方法名语法用途代码示例注意事项
StreamingContextnew StreamingContext(conf, batchDuration)创建流式上下文,指定配置和批处理间隔val conf = new SparkConf().setAppName("SocketWordCount")
val ssc = new StreamingContext(conf, Seconds(1))
批处理间隔设为 1 秒,影响处理频率。
socketTextStreamssc.socketTextStream(hostname, port)创建从 TCP socket 读取的 DStreamval lines = ssc.socketTextStream("localhost", 9999)hostname 和 port 需与发送端一致。
flatMapdstream.flatMap(line => line.split(" "))将每行文本拆分为单词序列val words = lines.flatMap(_.split(" "))返回的是 DStream[String]
mapdstream.map(word => (word, 1))将每个单词映射为键值对val wordPairs = words.map((_, 1))转换为 DStream[(String, Int)]
reduceByKeydstream.reduceByKey(_ + _)按 key 聚合值,统计词频val wordCounts = wordPairs.reduceByKey(_ + _)在批次内聚合,不跨批次。
printdstream.print()输出前 10 条结果到控制台wordCounts.print()仅用于调试,生产环境应写入外部系统。
startssc.start()启动流式计算,开始接收数据ssc.start()必须在定义完所有操作后调用。
awaitTerminationssc.awaitTermination()阻塞主线程,等待作业结束ssc.awaitTermination()若不调用,程序会立即退出。

完整代码逻辑流程(示意):

  1. 创建 SparkConf 和 StreamingContext
  2. 定义 socket DStream
  3. 执行 flatMap -> map -> reduceByKey -> print
  4. 启动上下文并等待终止

注意事项:

  • 运行前需在终端执行:nc -lk 9999 发送测试数据。
  • 程序运行后,每秒输出一次当前批次的词频统计。
  • 该示例不跨批次累计状态,仅统计当前批次数据。

第三章:DStream 基础操作

3.1 DStream 的创建方式

方法名语法用途代码示例注意事项
socketTextStreamssc.socketTextStream(hostname, port)从 TCP socket 读取文本流,常用于测试val lines = ssc.socketTextStream("localhost", 9999)仅用于开发调试,不具备容错性。
textFileStreamssc.textFileStream(directory)监控目录中新文件的写入,读取文本文件流val lines = ssc.textFileStream("hdfs://localhost:9000/logs/")文件必须原子性地移动到目录中,否则可能被重复处理。
fileStream[K, V, F]ssc.fileStreamK, V, F通用文件流,支持 SequenceFile、ObjectFile 等格式val stream = ssc.fileStreamLong, String, TextInputFormat需指定 Key、Value 类型和 InputFormat。
queueOfRDDsssc.queueStream(queueOfRDDs)将 RDD 队列转换为 DStream,用于模拟数据流val rddQueue = mutable.QueueRDD[Int]
val dstream = ssc.queueStream(rddQueue)
每个 RDD 视为一个批次,适合单元测试。
receiverStreamssc.receiverStream(myReceiver)使用自定义 Receiver 创建 DStreamval dstream = ssc.receiverStream(new CustomReceiver)自定义 Receiver 需继承 Receiver[T] 并实现 onStart/onStop
kafkaStream (Receiver)KafkaUtils.createStream(ssc, zkQuorum, group, topics)使用 Receiver 方式从 Kafka 消费数据KafkaUtils.createStream(ssc, "zk1:2181", "group1", Map("topic1" -> 1))需 ZooKeeper,Receiver 单点故障,不推荐生产使用。
kafkaDirectStream (Direct)KafkaUtils.createDirectStream(...)Direct 方式从 Kafka 拉取数据,无 ReceiverKafkaUtils.createDirectStream[String, String](ssc, kafkaParams, topics)推荐方式,支持精确一次语义,需手动管理 offset。

3.2 DStream 的转换操作(Transformations)

方法名语法用途代码示例注意事项
mapdstream.map(func)对 DStream 中每个元素应用函数val words = lines.map(_.split(" "))返回新 DStream,元素为转换结果。
flatMapdstream.flatMap(func)类似 map,但可返回多个元素(扁平化)val words = lines.flatMap(_.split(" "))常用于分词,结果展平为单个 DStream。
filterdstream.filter(func)过滤出满足条件的元素val errors = logs.filter(_.contains("ERROR"))返回布尔值的函数,保留 true 的记录。
uniondstream1.union(dstream2)合并两个同类型 DStreamval all = stream1.union(stream2)不去重,顺序不保证。
countdstream.count()返回每个批次中元素的个数(DStream[Long]val countStream = lines.count()统计批次内记录数。
countByValuedstream.countByValue()对元素按值计数,返回 (value, count) 流val counts = words.countByValue()适用于元素数量有限的场景,避免 OOM。
reducedstream.reduce(func)聚合批次内所有元素为单个值val sum = numbers.reduce(_ + _)函数需满足结合律,初始值隐式为第一个元素。
reduceByKeydstream.reduceByKey(func)按 key 聚合值,常用于 (K,V) 流val wordCounts = pairs.reduceByKey(_ + _)仅在当前批次内聚合,不跨批次。
groupByKeydstream.groupByKey()按 key 分组,value 聚合为 Iterableval grouped = pairs.groupByKey()可能导致大量数据 shuffle,慎用。
joindstream1.join(dstream2)对两个 (K,V) 和 (K,W) 流按 key 连接val joined = stream1.join(stream2)仅连接同一批次内的数据。
transformdstream.transform(func)对 DStream 内部的 RDD 应用任意函数stream.transform(rdd => rdd.join(otherRdd))强大但需谨慎使用,绕过 DStream 抽象。
updateStateByKeydstream.updateStateByKey(updateFunc)跨批次维护和更新每个 key 的状态val stateDstream = wordCounts.updateStateByKey(updateFunc)需启用 checkpoint,性能开销大。
mapWithStatedstream.mapWithState(stateSpec)高效的状态映射操作,推荐替代 updateStateByKeyval mapped = values.mapWithState(stateSpec)性能更好,支持状态超时,推荐使用。

3.3 DStream 的输出操作(Output Operations)

方法名语法用途代码示例注意事项
printdstream.print()输出前 10 个元素到控制台wordCounts.print()仅用于调试,生产环境避免使用。
foreachRDDdstream.foreachRDD(rdd => { ... })对每个批次的 RDD 执行自定义操作dstream.foreachRDD { rdd =>
if (!rdd.isEmpty) {
rdd.toDS().write.save("output/")
}
}
最常用的输出方式,可写入数据库、文件等。
saveAsTextFilesdstream.saveAsTextFiles(prefix, [suffix])将每个批次数据保存为文本文件wordCounts.saveAsTextFiles("output/prefix")文件路径包含时间戳,适合 HDFS。
saveAsObjectFilesdstream.saveAsObjectFiles(prefix)序列化 RDD 元素为对象文件stream.saveAsObjectFiles("data/output")需元素可序列化(Serializable)。
saveAsHadoopFilesdstream.saveAsHadoopFiles(prefix, ext, ...)使用 Hadoop OutputFormat 保存数据stream.saveAsHadoopFiles("out", "txt", classOf[Text], ...)支持自定义输出格式,如 SequenceFile。
saveAsNewAPIHadoopFilesdstream.saveAsNewAPIHadoopFiles(...)使用新的 Hadoop API 保存stream.saveAsNewAPIHadoopFiles("out", classOf[NullWritable], ...)适用于新版本 Hadoop。

foreachRDD 使用注意事项:

  • 每个 RDD 的操作应在 driver 端触发,避免在 executor 中创建 SparkContext。
  • 写入外部系统(如数据库)时,应使用 foreachPartition 或 connectionPool 避免频繁连接。

示例:写入 JDBC

dstream.foreachRDD { rdd =>
  rdd.foreachPartition { partition =>
    val conn = DriverManager.getConnection(url)
    partition.foreach { record =>
      val stmt = conn.prepareStatement("INSERT INTO table VALUES (?)")
      stmt.setString(1, record)
      stmt.executeUpdate()
    }
    conn.close()
  }
}

第四章:输入源(Input Sources)

4.1 基本数据源:Socket、文件流

方法名语法用途代码示例注意事项
socketTextStreamssc.socketTextStream(host, port)从 TCP socket 读取 UTF-8 文本val lines = ssc.socketTextStream("localhost", 9999)无容错,Receiver 失败则数据丢失;仅用于测试。
textFileStreamssc.textFileStream(directory)监控目录中新文件,读取文本内容val logs = ssc.textFileStream("file:///var/logs/")文件必须通过”移动”方式写入目录,否则可能被重复读取。
binaryRecordsStreamssc.binaryRecordsStream(path, recordLength)读取固定长度的二进制记录文件val binary = ssc.binaryRecordsStream("data/", 1024)适用于日志、传感器等二进制数据。

4.2 高级数据源:Kafka(Receiver 与 Direct 方式)

Receiver 方式(已过时,不推荐)

方法名语法用途代码示例注意事项
createStreamKafkaUtils.createStream(ssc, zkQuorum, groupID, topics)使用 Receiver 从 Kafka 消费数据val stream = KafkaUtils.createStream(ssc, "zk:2181", "group1", Map("topic" -> 1))使用 ZooKeeper 管理 offset;Receiver 单点故障;不支持精确一次。

Direct 方式(推荐)

方法名语法用途代码示例注意事项
createDirectStreamKafkaUtils.createDirectStream(ssc, kafkaParams, topics)Direct 方式拉取 Kafka 数据,无 Receiverval kafkaParams = Map("bootstrap.servers" -> "kafka:9092",
"group.id" -> "group1",
"key.deserializer" -> classOf[StringDeserializer],
"value.deserializer" -> classOf[StringDeserializer])
val stream = KafkaUtils.createDirectStream[String, String](ssc, kafkaParams, Set("topic1"))
每个批次直接调用 Kafka 消费者 API;支持精确一次语义;offset 可手动提交。
获取 offsetstream.foreachRDD { rdd =>
val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
// 处理逻辑
// 可在此提交 offset 到 Kafka 或外部存储
}
获取当前批次的 offset 范围见上需将 RDD 转换为 HasOffsetRanges 类型。

Direct 方式优势:

  • 无单点故障
  • 更好的容错性
  • 支持从任意 offset 开始消费
  • 可以与 checkpoint 结合实现精确一次处理

4.3 Flume 与 Kinesis 集成

Flume 集成

方法名语法用途代码示例注意事项
createStreamFlumeUtils.createStream(ssc, host, port)从 Flume agent 接收数据(Push-based)val stream = FlumeUtils.createStream(ssc, "flume-agent", 41414)需 Flume agent 配置 Avro Source。
createPollingStreamFlumeUtils.createPollingStream(ssc, hosts)轮询方式从多个 Flume agent 拉取数据val stream = FlumeUtils.createPollingStream(ssc, List("agent1", "agent2"))更可靠,支持负载均衡。

Kinesis 集成(AWS)

方法名语法用途代码示例注意事项
createStreamKinesisUtils.createStream(ssc, appName, streamName, endpointUrl, ...)从 AWS Kinesis 流读取数据val stream = KinesisUtils.createStream(ssc, "myApp", "myStream", "https://kinesis.us-east-1.amazonaws.com", awsAccessKey, awsSecretKey, InitialPositionInStream.LATEST, 4)需 AWS 凭证;适用于云环境数据采集。
createDirectStreamKinesisUtils.createDirectStream(...)推荐方式,Direct 模式读取 Kinesis推荐使用 Kinesis Client Library (KCL) 集成性能更好,支持 checkpoint 到 DynamoDB。

4.4 自定义数据源

概念/方法说明注意事项
继承 Receiver[T]自定义数据源需继承 org.apache.spark.streaming.receiver.Receiver[T] 抽象类必须实现 onStart()onStop() 方法。
onStart()启动数据接收逻辑,通常在新线程中拉取数据可启动线程池、网络连接等;调用 store(data) 保存数据。
onStop()停止数据接收,释放资源关闭连接、中断线程,确保优雅退出。

store(data) - 将接收到的数据存入 Spark

def onStart(): Unit = {
  new Thread("Custom Receiver") {
    override def run(): Unit = {
      while (!isStopped()) {
        val data = fetchData()
        store(data)
      }
    }
  }.start()
}
概念/方法说明注意事项
容错性Receiver 重启后可能重复接收数据建议在数据源端支持幂等性或 offset 记录。
性能Receiver 运行在 Executor 中,占用一个核心建议使用 Direct 方式替代,避免资源竞争。

第五章:窗口操作(Windowed Operations)

5.1 窗口、滑动间隔与批处理间隔

概念说明注意事项
批处理间隔(Batch Interval)数据流被切分的最小时间单位(如 1s),每个批次生成一个 RDD。决定系统调度频率,影响延迟与吞吐量。
窗口长度(Window Length)窗口操作覆盖的时间范围(如 30s),表示参与计算的数据时间跨度。必须是批处理间隔的整数倍。
滑动间隔(Slide Interval)窗口滑动的时间步长(如 10s),决定多久计算一次窗口结果。也必须是批处理间隔的整数倍;若等于窗口长度,则为”滚动窗口”。
窗口操作原理将多个连续批次的数据合并,在指定窗口长度内进行聚合或转换,每隔滑动间隔执行一次。例如:每 10 秒统计过去 30 秒的词频。
示例配置batchDuration = 5s, windowLength = 20s, slideInterval = 10s → 每 10s 计算最近 20s 数据系统会自动组合 4 个批次(20/5)为一个窗口。

5.2 常用窗口转换操作

方法名语法用途代码示例注意事项
windowdstream.window(windowLength, slideInterval)返回一个基于窗口的 DStream,后续可进行任意转换val windowedStream = stream.window(Seconds(30), Seconds(10))返回新的 DStream,可链式调用 map、reduce 等操作。
countByWindowdstream.countByWindow(windowLength, slideInterval)统计窗口内元素总数val countStream = stream.countByWindow(Seconds(60), Seconds(10))高效实现,避免手动 reduce。
reduceByWindowdstream.reduceByWindow(reduceFunc, invReduceFunc, windowLength, slideInterval)增量式聚合窗口数据,支持反向减少(高效)val sumStream = nums.reduceByWindow(_ + _, _ - _, Seconds(30), Seconds(10))提供 invReduceFunc 可实现 O(1) 滑动更新,否则为 O(N)。
reduceByWindow(无反向函数)dstream.reduceByWindow(reduceFunc, windowLength, slideInterval)普通窗口聚合,不支持增量更新val sum = nums.reduceByWindow(_ + _, Seconds(30), Seconds(10))每次重新计算整个窗口,性能较差。
countByValueAndWindowdstream.countByValueAndWindow(windowLength, slideInterval)按值统计窗口内各元素出现次数val counts = words.countByValueAndWindow(Seconds(60), Seconds(10))类似 countByValue,但作用于窗口数据。
union + windowstream1.union(stream2).window(...)合并多个流后进行窗口操作val combined = src1.union(src2).window(Seconds(20), Seconds(5))可用于多源数据的统一窗口分析。

reduceByWindow 增量更新说明:

  • reduceFunc:正向聚合函数(如 _ + _
  • invReduceFunc:反向减少函数(如 _ - _),用于移除滑出窗口的数据
  • 若提供反向函数,Spark 可基于前一个窗口结果增量更新,极大提升性能

5.3 基于窗口的状态管理

概念/方法说明注意事项
窗口状态窗口操作本身不保存跨窗口的状态,每次计算独立若需跨多个窗口累积状态,需结合 updateStateByKeymapWithState
状态持久化窗口计算结果可通过 foreachRDD 写入外部存储(如 Redis、数据库)实现持久化避免在 driver 内存中累积大量状态,防止 OOM。
与 mapWithState 结合可将窗口聚合结果作为输入,传递给 mapWithState 进行长期状态维护例如:每 10 秒统计过去 1 分钟的 UV,再累计全天 PV。
Checkpoint 支持窗口操作本身不依赖 checkpoint,但若使用有状态操作(如 updateStateByKey),需启用 checkpoint必须设置 checkpoint 目录以支持故障恢复。

第六章:有状态流处理

6.1 updateStateByKey 操作

方法名语法用途代码示例注意事项
updateStateByKeydstream.updateStateByKey(updateFunc)跨批次维护每个 key 的状态,常用于累计计数def updateFunc(values: Seq[Int], state: Option[Int]): Option[Int] = {
val currentCount = values.sum
val previousCount = state.getOrElse(0)
Some(currentCount + previousCount)
}
val stateDStream = wordCounts.updateStateByKey(updateFunc)
需启用 checkpoint,性能开销大,不推荐用于高吞吐场景。
updateFunc 参数(Seq[T], Option[S]) => Option[S]输入为当前批次该 key 的值序列和之前状态见上必须返回新的状态(Some)或删除状态(None)。
Checkpoint 依赖必须设置 ssc.checkpoint(checkpointDir)ssc.checkpoint("hdfs://localhost:9000/checkpoint")否则运行会报错。
性能问题每个批次处理所有 key,包括无更新的 key导致大量无效计算,吞吐量低。

6.2 mapWithState 操作(推荐方式)

方法名语法用途代码示例注意事项
mapWithStatedstream.mapWithState(stateSpec)高效的有状态映射操作,仅处理有数据输入的 keyval spec = StateSpec.function(updateFunc).timeout(Minutes(10))
val mappedDStream = keyValueStream.mapWithState(stateSpec)
推荐替代 updateStateByKey,性能更好。
StateSpec.functionStateSpec.function(updateFunc)定义状态更新函数def updateFunc(key: String, value: Option[Int], state: State[Int]): Option[(String, Int)] = {
val sum = value.getOrElse(0) + state.getOption().getOrElse(0)
state.update(sum)
Some((key, sum))
}
函数返回输出数据流元素。
State 对象方法state.get(), state.update(), state.exists(), state.remove()用于读取、更新、判断、删除状态见上支持状态超时(TTL)机制。
timeoutspec.timeout(duration)设置状态空闲超时时间,超时后自动清除spec.timeout(Minutes(5))减少内存占用,避免状态无限增长。
Checkpoint同样需要启用 checkpointssc.checkpoint("hdfs://...")用于故障恢复。

mapWithState 优势:

  • 仅处理有输入的 key,跳过无更新的 key
  • 支持状态超时(TTL)
  • 可返回输出流,灵活性高
  • 性能显著优于 updateStateByKey

6.3 状态管理的最佳实践

实践建议说明注意事项
使用 mapWithState 替代 updateStateByKey性能更好,支持 TTL,推荐用于新项目老项目可逐步迁移。
启用 Checkpoint必须设置 checkpoint 目录以支持状态恢复建议使用 HDFS 等可靠存储。
控制状态大小避免存储过多 key 或大对象可通过 timeout 清理过期状态。
外部状态存储对于超大状态,可将状态存入 Redis、RocksDB 等Spark Streaming 本身不支持嵌入式状态后端(如 Flink 的 StateBackend)。
监控状态增长通过 Spark UI 查看任务处理时间和内存使用防止 OOM 导致 executor 崩溃。
精确一次语义结合 Direct Kafka + 手动提交 offset + 幂等写入实现端到端精确一次处理。

第七章:容错与可靠性

7.1 Spark Streaming 的容错机制

机制说明注意事项
RDD Lineage每个 DStream 对应 RDD 序列,通过血统(lineage)可重新计算丢失的批次基于 Spark Core 的容错能力,无需复制数据。
批次确认Receiver 接收数据后,需确认已复制到多个 executor 才向外部系统确认确保数据不丢失(Receiver 模式)。
Checkpoint定期保存元数据(配置、DStream 操作、定时器)和状态到可靠存储用于 driver 故障恢复。
Write Ahead Log (WAL)Receiver 模式下,接收到的数据先写入 WAL,再处理防止 Receiver 失败导致数据丢失,但增加延迟。

7.2 数据接收的可靠性(Receiver 与 Direct 模式对比)

对比项Receiver 模式Direct 模式(Kafka)
架构使用 Receiver 接收数据,数据存储在 executor 内存直接调用 Kafka 消费者 API 拉取数据
容错性依赖 WAL 和复制确保不丢失更可靠,每个批次精确控制 offset
语义支持最多一次(无 WAL)或至少一次(有 WAL)可实现精确一次(结合幂等写入)
性能存在单点瓶颈,Receiver 可能成为瓶颈并行度高,每个 partition 一个 task
Offset 管理由 ZooKeeper 管理由 Spark 自行管理,可手动提交
推荐程度已过时,不推荐强烈推荐用于生产环境

7.3 检查点(Checkpointing)机制

方法/概念语法/说明用途注意事项
checkpointssc.checkpoint(directory)设置 checkpoint 目录ssc.checkpoint("hdfs://localhost:9000/checkpoint")
Checkpoint 内容• 应用配置
• DStream 操作图
• 批次时间信息
• 有状态操作的元数据(如 updateStateByKey)
用于 driver 重启后恢复执行状态不保存 RDD 数据本身。
触发频率默认每 10 个批次 checkpoint 一次可通过 ssc.checkpoint() 和 Duration 控制频繁 checkpoint 影响性能。
目录要求必须为 HDFS、S3 等容错文件系统本地文件系统不支持高可用单机测试可使用本地路径,但生产环境禁用。

7.4 故障恢复流程

故障类型恢复机制说明注意事项
Executor 故障RDD Lineage 重新计算丢失的批次可通过 lineage 重新生成数据仍在内存或 WAL 中时可恢复。
Driver 故障Checkpoint 恢复从 checkpoint 目录重建 StreamingContextval ssc = StreamingContext.getOrCreate(checkpointPath, createStreamingContext)
Kafka Offset 丢失Direct 模式手动管理从 checkpoint 或外部存储读取上一次 offset确保写入外部系统(如 DB)与 offset 提交原子性。
数据重复至少一次语义导致通过幂等操作或外部去重避免影响如写入 Redis 使用 SET 命令,或数据库使用唯一键。

getOrCreate 使用示例:

def createStreamingContext(): StreamingContext = {
  val ssc = new StreamingContext(...)
  ssc.checkpoint("hdfs://...")
  // 定义 DStream 操作
  ssc
}
val ssc = StreamingContext.getOrCreate("hdfs://checkpoint", createStreamingContext)
ssc.start()
ssc.awaitTermination()

第八章:性能调优

8.1 批处理间隔(Batch Duration)设置

概念/方法说明注意事项
批处理间隔(Batch Duration)将数据流切分为微批的时间间隔(如 500ms、1s、5s),是性能调优的核心参数间隔越短,延迟越低,但调度开销越大;需在延迟与吞吐间权衡。
设置方式在创建 StreamingContext 时指定val ssc = new StreamingContext(sparkConf, Seconds(1))
调优建议• 初始设置为 1-2 秒
• 观察单个批次处理时间(Processing Time)
• 确保处理时间 < 批处理间隔,避免积压
若处理时间 > 批处理间隔,系统将无法实时处理,导致延迟累积。
监控指标通过 Spark UI 查看每个批次的:
• 调度延迟(Scheduling Delay)
• 处理时间(Processing Time)
调度延迟应尽可能小,理想为 0。

8.2 并行度与任务调度优化

方法/配置语法/说明用途注意事项
增加 Kafka 分区数提高 Kafka topic 的 partition 数量增加 Spark 并行消费 task 数Spark 中每个 partition 对应一个 task,分区数决定最大并行度。
spark.default.parallelism--conf "spark.default.parallelism=4"设置默认并行度,影响 shuffle 操作的 task 数建议设为集群 CPU 核心总数。
spark.streaming.concurrentJobs--conf "spark.streaming.concurrentJobs=2"允许同时处理多个 streaming 作业默认为 1;若作业间无依赖,可提高并发。
减少任务启动开销使用 foreachPartition 而非 foreach批量处理 partition 数据,减少外部连接创建避免在 foreach 中频繁创建数据库连接。
数据倾斜处理使用 mapPartitions 手动重分区或加盐(salting)避免某些 task 处理过多数据可通过 repartitioncoalesce 调整分区数。

8.3 数据序列化与内存管理

配置项语法用途注意事项
spark.serializerspark.serializer org.apache.spark.serializer.KryoSerializer使用 Kryo 序列化,比 Java 序列化更高效必须注册自定义类以获得最佳性能。
spark.kryo.classesToRegisterspark.kryo.classesToRegister com.example.Event注册需要 Kryo 序列化的类提高序列化速度,减少内存占用。
spark.streaming.kafka.maxRatePerPartitionkafkaParams.put("maxRatePerPartition", "1000")限制每个 Kafka 分区每秒拉取的消息数防止 executor 内存溢出(OOM),尤其在流量突增时。
spark.streaming.backpressure.enabledspark.streaming.backpressure.enabled true启用背压机制,自动调节拉取速率需关闭 maxRatePerPartition
内存分配spark.executor.memory=4g合理分配 executor 内存留出足够空间给 stream buffering 和 shuffle。

8.4 背压机制(Backpressure)启用

方法/配置语法用途注意事项
启用背压sparkConf.set("spark.streaming.backpressure.enabled", "true")系统根据处理能力自动调节数据摄入速率无需设置 maxRatePerPartition,否则会覆盖背压。
背压初始速率spark.streaming.backpressure.initialRate设置背压机制的初始摄入速率(如 1000 msg/s)可根据历史流量设置合理初始值。
工作原理根据前一批次的处理时间动态调整下一批次的拉取量处理慢则减少摄入,处理快则增加摄入实现自动负载均衡,防止系统过载。
监控指标通过 Spark UI 查看 RateEstimator 的 estimatedRate观察系统实时摄入速率变化背压生效时,速率会随负载波动。

8.5 Kafka 消费性能调优

配置项语法/说明用途注意事项
fetch.message.max.byteskafkaParams.put("fetch.message.max.bytes", "10485880")单次拉取最大数据量(默认 1MB)可适当增大以提高吞吐,但需匹配 executor 内存。
max.partition.fetch.byteskafkaParams.put("max.partition.fetch.bytes", "10485880")每个分区单次拉取最大字节数建议与 fetch.message.max.bytes 一致。
auto.offset.resetkafkaParams.put("auto.offset.reset", "latest")无初始 offset 时从最新(latest)或最早(earliest)开始生产环境建议设为 latest 避免重放历史数据。
手动提交 offsetforeachRDD 中调用 commitAsync实现精确一次语义需确保写入外部系统与 offset 提交在同一个语义下。
并行度匹配Kafka 分区数 ≥ Spark 分区数确保最大并行消费能力若 Kafka 分区少,Spark 无法并行处理。

第九章:与外部系统集成

9.1 写入 Kafka

方法/库说明代码示例注意事项
使用 foreachRDD(推荐)通过 foreachRDD 写入 Kafka,使用 KafkaProducerdstream.foreachRDD { rdd =>
rdd.foreachPartition { partition =>
val producer = new KafkaProducer[String, String]()
partition.foreach { record =>
producer.send(new ProducerRecord[String, String]("topic", value))
}
producer.close()
}
}
使用连接池或单例 producer 可提升性能。
幂等生产者enable.idempotence=true防止消息重复Kafka 0.11+ 支持,推荐开启。
事务写入producer.initTransactions() + send + commitTransaction实现精确一次语义需 Kafka 配置支持事务。

9.2 写入数据库(JDBC、MySQL、PostgreSQL)

方法说明代码示例注意事项
foreachPartition + 批量插入每个 partition 建立连接,批量执行 SQLdstream.foreachRDD { rdd =>
rdd.foreachPartition { partition =>
val conn = DriverManager.getConnection(url)
val stmt = conn.prepareStatement("INSERT INTO table VALUES (?, ?)")
partition.foreach { case (k, v) =>
stmt.setString(1, k)
stmt.setString(2, v)
stmt.addBatch()
}
stmt.executeBatch()
conn.close()
}
}
使用 addBatch + executeBatch 提升性能。
连接池使用 HikariCP、DBCP 等连接池复用数据库连接,避免频繁创建foreachPartition 外部初始化连接池。
幂等设计使用唯一约束或 INSERT IGNORE / ON CONFLICT防止重复数据插入MySQL 用 INSERT IGNORE,PostgreSQL 用 ON CONFLICT DO NOTHING
批量大小控制限制每批插入记录数(如 1000 条)防止事务过大导致锁表或超时可通过 take + 循环实现分批。

9.3 写入 Redis

方法说明代码示例注意事项
foreachPartition + Jedis每个 partition 使用 Jedis 客户端写入dstream.foreachRDD { rdd =>
rdd.foreachPartition { partition =>
val jedis = new Jedis("localhost")
partition.foreach { case (k, v) =>
jedis.set(k, v)
}
jedis.close()
}
}
推荐使用连接池(如 JedisPool)。
批量操作使用 pipelined 提交多个命令val pipeline = jedis.pipelined()
partition.foreach { case (k, v) => pipeline.set(k, v) }
pipeline.sync()
显著提升写入吞吐量。
数据结构选择根据需求选择 String、Hash、Set 等如用户画像可用 Hash 存储合理设计 key 结构,避免 key 过多。
超时设置设置 key 过期时间(TTL)jedis.setex("key", 3600, "value")防止 Redis 内存无限增长。

9.4 写入 HDFS / Parquet / Delta Lake

格式方法代码示例注意事项
文本文件(HDFS)dstream.saveAsTextFiles("hdfs://path/prefix")lines.saveAsTextFiles("hdfs://namenode:9000/output/logs")文件名包含时间戳,适合日志存储。
Parquet(结构化)转换为 DataFrame 后写入dstream.foreachRDD { rdd =>
if (!rdd.isEmpty()) {
val df = rdd.toDS().toDF()
df.write.mode("append").parquet("hdfs://path/parquet")
}
}
需引入 spark-sql 模块;Parquet 支持列式存储、压缩。
Delta Lake使用 Delta Lake 表格式df.write.format("delta").mode("append").save("/path/delta-table")支持 ACID、时间旅行、Schema 治理;需引入 Delta Lake 依赖。
小文件问题合并小文件或设置 checkpoint使用 coalesce 减少输出分区数df.coalesce(1).write... 可减少文件数,但降低并行度。
分区写入按时间或 key 分区存储df.write.partitionBy("date").parquet("path")便于后续查询过滤。

Delta Lake 优势:

  • 支持流批统一
  • 提供精确一次语义
  • 兼容 Spark SQL
  • 推荐用于生产环境数据湖构建

第十章:结构化流式处理(Structured Streaming)简介

10.1 Structured Streaming 核心概念

概念说明注意事项
数据流作为”无限表”(Unbounded Table)将流数据视为持续追加的动态表,批处理与流处理统一模型基于 Spark SQL 引擎,支持 DataFrame/Dataset API。
流式查询(Streaming Query)使用 writeStream 启动持续查询,结果不断更新查询可运行在 Append、Update、Complete 输出模式下。
触发器(Triggers)控制查询执行频率,如微批(默认)或一次性处理支持 ProcessingTime、Once、Continuous 模式。
事件时间(Event Time)基于数据本身的时间戳而非处理时间进行计算支持窗口聚合、延迟数据处理。
水印(Watermark)定义可容忍的延迟时间,用于清理状态防止状态无限增长,需合理设置(如 withWatermark("eventTime", "10 minutes"))。
端到端精确一次语义结合幂等写入与事务提交,确保结果只处理一次依赖 Source 和 Sink 支持(如 Kafka + Delta Lake)。

10.2 与 DStream API 对比

对比项DStream APIStructured Streaming
编程模型基于 RDD 的离散流(DStream)基于 DataFrame/Dataset 的结构化流
API 风格函数式(map, reduce, window)SQL/DSL 风格,支持强类型 Dataset
容错机制依赖 Checkpoint 和 WAL自动管理,无需手动设置 checkpoint
语义保证至少一次或手动实现精确一次默认支持端到端精确一次(Source + Sink 支持时)
状态管理updateStateByKey / mapWithState自动管理,通过 groupByKey + agg + watermark
时间处理仅处理时间(Processing Time)支持事件时间(Event Time)和水印
输出模式foreachRDD 自定义输出支持 Append、Update、Complete 三种模式
性能优化手动调优(并行度、序列化等)Catalyst 优化器自动优化执行计划
推荐程度已过时,维护模式强烈推荐用于新项目

10.3 使用 Structured Streaming 实现 WordCount

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._

val spark = SparkSession.builder()
  .appName("StructuredWordCount")
  .config("spark.sql.streaming.checkpointLocation", "/checkpoint/dir")
  .getOrCreate()

import spark.implicits._

// 1. 从 Socket 读取数据(仅用于测试)
val lines = spark.readStream
  .format("socket")
  .option("host", "localhost")
  .option("port", 9999)
  .load()

// 2. 分词并计数
val words = lines.as[String].flatMap(_.split(" "))
val wordCounts = words.groupBy("value").count()

// 3. 启动流式查询
val query = wordCounts.writeStream
  .outputMode("complete")  // 因为是聚合,使用 complete 模式
  .format("console")
  .trigger(Trigger.ProcessingTime("5 seconds"))
  .start()

query.awaitTermination()

说明:

  • outputMode("complete"):输出所有聚合结果(适合小数据量)
  • trigger:每 5 秒执行一次微批处理
  • checkpointLocation:必须设置,用于故障恢复

10.4 流式聚合与事件时间处理

操作语法用途注意事项
事件时间窗口.groupBy(window($"eventTime", "10 minutes", "5 minutes"), $"userId").count()基于事件时间进行滑动窗口聚合必须启用水印。
添加水印.withWatermark("eventTime", "10 minutes")允许延迟 10 分钟的数据,超时则丢弃水印必须在聚合前设置。
会话窗口.groupBy(session_window($"eventTime", "30 minutes")).count()处理用户会话(无活动超时则结束)Spark 3.0+ 支持。
去重.dropDuplicates(Seq("userId", "eventId"), $"eventTime")基于事件时间和 key 去重需水印支持,保留每个 key 的最新记录。
状态清理水印机制自动清理过期状态防止状态无限增长水印时间越长,状态保留越久。

事件时间处理示例:

val withEventTime = df.select(
  $"value".cast("string"),
  $"timestamp".cast("timestamp").as("eventTime")
)

val withWatermark = withEventTime.withWatermark("eventTime", "5 minutes")

val windowedCounts = withWatermark
  .groupBy(window($"eventTime", "10 minutes"), $"value")
  .count()

第十一章:实际项目案例

11.1 实时日志分析系统

组件技术选型说明
数据源Flume / Logstash / Kafka收集 Nginx、应用日志并写入 Kafka
流处理引擎Structured Streaming从 Kafka 读取日志,解析为结构化数据
处理逻辑• 正则解析日志
• 过滤错误日志(ERROR/500)
• 统计每分钟请求数、错误率
支持事件时间窗口聚合
存储• Redis:实时错误计数
• HDFS/Parquet:原始日志归档
• Elasticsearch:日志检索
支持多系统写入
可视化Kibana / Grafana实时监控仪表盘

11.2 实时广告点击流处理

组件技术选型说明
数据源Kafka(点击日志)每条记录包含 userId、adId、timestamp、type(click/impression)
流处理引擎Structured Streaming实时计算 CTR(点击率)
处理逻辑• 分别统计点击和曝光数量
• 每 10 秒计算 CTR = clicks / impressions
• 支持按广告位、用户群体分组
使用事件时间 + 水印处理延迟数据
存储• Redis:实时 CTR 缓存
• Delta Lake:持久化分析结果
支持 A/B 测试分析
报警Prometheus + AlertmanagerCTR 异常下降时触发告警

11.3 用户行为实时监控

组件技术选型说明
数据源Kafka(埋点日志)包含页面浏览、按钮点击、停留时长等事件
流处理引擎Structured Streaming实时识别异常行为(如频繁登录失败)
处理逻辑• 会话划分(session_window)
• 用户行为序列分析
• 实时 UV/PV 统计
• 异常检测(规则引擎或模型)
支持低延迟响应
存储• Redis:实时用户画像
• HBase:行为历史存储
• Kafka:转发告警事件
支持实时推荐与风控
应用• 实时推荐系统
• 风控系统(防刷)
• 用户留存分析
业务驱动实时处理

第十二章:常见问题与调试

12.1 延迟问题排查

现象可能原因解决方案
调度延迟(Scheduling Delay)高批处理间隔过短或资源不足增加批处理间隔或集群资源
处理时间(Processing Time)长数据倾斜、shuffle 慢、外部 I/O 瓶颈优化 SQL、增加并行度、使用连接池
端到端延迟高数据源生产慢或网络延迟检查 Kafka 消费 lag、网络状况
背压未生效maxRatePerPartition 设置过高或背压未启用关闭限速,启用 backpressure.enabled

12.2 数据丢失与重复问题

问题原因解决方案
数据丢失• Receiver 模式未启用 WAL
• Kafka 消费 offset 提交过早
使用 Direct 模式,确保数据处理后再提交 offset
数据重复• 至少一次语义
• 外部系统写入失败重试
实现幂等写入(如 Redis SET、数据库唯一键)
Kafka 重复消费offset 提交与处理不同步使用事务或两阶段提交
精确一次实现结合幂等 Sink 与事务提交Kafka Producer + Delta Lake Sink 支持

12.3 检查点异常处理

异常原因解决方案
Checkpoint directory not set未调用 ssc.checkpoint()添加 ssc.checkpoint("/path")
Checkpoint 写入失败HDFS 权限不足或磁盘满检查存储系统状态与权限
恢复失败Checkpoint 数据损坏或版本不兼容清除 checkpoint 目录重新开始(有状态丢失风险)
性能下降Checkpoint 频繁调整 checkpointDuration,减少频率

12.4 Web UI 监控指标解读

页面关键指标说明
Streaming 页• Batch Processing Time:处理时间应 < 批处理间隔
• Scheduling Delay:理想为 0
• Total Delay:处理时间 + 调度延迟
延迟过高表示系统过载
Jobs / Stages• Task 执行时间分布
• Shuffle 读写量
• GC 时间
识别数据倾斜或资源瓶颈
Executors• 内存使用率
• 磁盘溢出(Spill)
• 线程阻塞
内存不足时增加 executor.memory
Kafka SourceoffsetLog 中的 currentOffset vs endOffset差值大表示消费滞后(lag)

建议: 定期监控 UI,设置 Prometheus + Grafana 做长期趋势分析。