Article
第一章: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 Streaming | Spark 的流式计算模块,用于构建可扩展、容错的实时数据处理应用。它将流数据划分为一系列小的批处理任务,在 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_HOME 和 PATH。 | 本地模式无需集群,但生产环境需部署 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.version | Scala 版本 | 2.12.15 |
| spark.version | Spark 版本 | 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 示例)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| StreamingContext | new StreamingContext(conf, batchDuration) | 创建流式上下文,指定配置和批处理间隔 | val conf = new SparkConf().setAppName("SocketWordCount")val ssc = new StreamingContext(conf, Seconds(1)) | 批处理间隔设为 1 秒,影响处理频率。 |
| socketTextStream | ssc.socketTextStream(hostname, port) | 创建从 TCP socket 读取的 DStream | val lines = ssc.socketTextStream("localhost", 9999) | hostname 和 port 需与发送端一致。 |
| flatMap | dstream.flatMap(line => line.split(" ")) | 将每行文本拆分为单词序列 | val words = lines.flatMap(_.split(" ")) | 返回的是 DStream[String]。 |
| map | dstream.map(word => (word, 1)) | 将每个单词映射为键值对 | val wordPairs = words.map((_, 1)) | 转换为 DStream[(String, Int)]。 |
| reduceByKey | dstream.reduceByKey(_ + _) | 按 key 聚合值,统计词频 | val wordCounts = wordPairs.reduceByKey(_ + _) | 在批次内聚合,不跨批次。 |
dstream.print() | 输出前 10 条结果到控制台 | wordCounts.print() | 仅用于调试,生产环境应写入外部系统。 | |
| start | ssc.start() | 启动流式计算,开始接收数据 | ssc.start() | 必须在定义完所有操作后调用。 |
| awaitTermination | ssc.awaitTermination() | 阻塞主线程,等待作业结束 | ssc.awaitTermination() | 若不调用,程序会立即退出。 |
完整代码逻辑流程(示意):
- 创建 SparkConf 和 StreamingContext
- 定义 socket DStream
- 执行
flatMap -> map -> reduceByKey -> print - 启动上下文并等待终止
注意事项:
- 运行前需在终端执行:
nc -lk 9999发送测试数据。 - 程序运行后,每秒输出一次当前批次的词频统计。
- 该示例不跨批次累计状态,仅统计当前批次数据。
第三章:DStream 基础操作
3.1 DStream 的创建方式
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| socketTextStream | ssc.socketTextStream(hostname, port) | 从 TCP socket 读取文本流,常用于测试 | val lines = ssc.socketTextStream("localhost", 9999) | 仅用于开发调试,不具备容错性。 |
| textFileStream | ssc.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。 |
| queueOfRDDs | ssc.queueStream(queueOfRDDs) | 将 RDD 队列转换为 DStream,用于模拟数据流 | val rddQueue = mutable.QueueRDD[Int]val dstream = ssc.queueStream(rddQueue) | 每个 RDD 视为一个批次,适合单元测试。 |
| receiverStream | ssc.receiverStream(myReceiver) | 使用自定义 Receiver 创建 DStream | val 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 拉取数据,无 Receiver | KafkaUtils.createDirectStream[String, String](ssc, kafkaParams, topics) | 推荐方式,支持精确一次语义,需手动管理 offset。 |
3.2 DStream 的转换操作(Transformations)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| map | dstream.map(func) | 对 DStream 中每个元素应用函数 | val words = lines.map(_.split(" ")) | 返回新 DStream,元素为转换结果。 |
| flatMap | dstream.flatMap(func) | 类似 map,但可返回多个元素(扁平化) | val words = lines.flatMap(_.split(" ")) | 常用于分词,结果展平为单个 DStream。 |
| filter | dstream.filter(func) | 过滤出满足条件的元素 | val errors = logs.filter(_.contains("ERROR")) | 返回布尔值的函数,保留 true 的记录。 |
| union | dstream1.union(dstream2) | 合并两个同类型 DStream | val all = stream1.union(stream2) | 不去重,顺序不保证。 |
| count | dstream.count() | 返回每个批次中元素的个数(DStream[Long]) | val countStream = lines.count() | 统计批次内记录数。 |
| countByValue | dstream.countByValue() | 对元素按值计数,返回 (value, count) 流 | val counts = words.countByValue() | 适用于元素数量有限的场景,避免 OOM。 |
| reduce | dstream.reduce(func) | 聚合批次内所有元素为单个值 | val sum = numbers.reduce(_ + _) | 函数需满足结合律,初始值隐式为第一个元素。 |
| reduceByKey | dstream.reduceByKey(func) | 按 key 聚合值,常用于 (K,V) 流 | val wordCounts = pairs.reduceByKey(_ + _) | 仅在当前批次内聚合,不跨批次。 |
| groupByKey | dstream.groupByKey() | 按 key 分组,value 聚合为 Iterable | val grouped = pairs.groupByKey() | 可能导致大量数据 shuffle,慎用。 |
| join | dstream1.join(dstream2) | 对两个 (K,V) 和 (K,W) 流按 key 连接 | val joined = stream1.join(stream2) | 仅连接同一批次内的数据。 |
| transform | dstream.transform(func) | 对 DStream 内部的 RDD 应用任意函数 | stream.transform(rdd => rdd.join(otherRdd)) | 强大但需谨慎使用,绕过 DStream 抽象。 |
| updateStateByKey | dstream.updateStateByKey(updateFunc) | 跨批次维护和更新每个 key 的状态 | val stateDstream = wordCounts.updateStateByKey(updateFunc) | 需启用 checkpoint,性能开销大。 |
| mapWithState | dstream.mapWithState(stateSpec) | 高效的状态映射操作,推荐替代 updateStateByKey | val mapped = values.mapWithState(stateSpec) | 性能更好,支持状态超时,推荐使用。 |
3.3 DStream 的输出操作(Output Operations)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
dstream.print() | 输出前 10 个元素到控制台 | wordCounts.print() | 仅用于调试,生产环境避免使用。 | |
| foreachRDD | dstream.foreachRDD(rdd => { ... }) | 对每个批次的 RDD 执行自定义操作 | dstream.foreachRDD { rdd => if (!rdd.isEmpty) { rdd.toDS().write.save("output/") }} | 最常用的输出方式,可写入数据库、文件等。 |
| saveAsTextFiles | dstream.saveAsTextFiles(prefix, [suffix]) | 将每个批次数据保存为文本文件 | wordCounts.saveAsTextFiles("output/prefix") | 文件路径包含时间戳,适合 HDFS。 |
| saveAsObjectFiles | dstream.saveAsObjectFiles(prefix) | 序列化 RDD 元素为对象文件 | stream.saveAsObjectFiles("data/output") | 需元素可序列化(Serializable)。 |
| saveAsHadoopFiles | dstream.saveAsHadoopFiles(prefix, ext, ...) | 使用 Hadoop OutputFormat 保存数据 | stream.saveAsHadoopFiles("out", "txt", classOf[Text], ...) | 支持自定义输出格式,如 SequenceFile。 |
| saveAsNewAPIHadoopFiles | dstream.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、文件流
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| socketTextStream | ssc.socketTextStream(host, port) | 从 TCP socket 读取 UTF-8 文本 | val lines = ssc.socketTextStream("localhost", 9999) | 无容错,Receiver 失败则数据丢失;仅用于测试。 |
| textFileStream | ssc.textFileStream(directory) | 监控目录中新文件,读取文本内容 | val logs = ssc.textFileStream("file:///var/logs/") | 文件必须通过”移动”方式写入目录,否则可能被重复读取。 |
| binaryRecordsStream | ssc.binaryRecordsStream(path, recordLength) | 读取固定长度的二进制记录文件 | val binary = ssc.binaryRecordsStream("data/", 1024) | 适用于日志、传感器等二进制数据。 |
4.2 高级数据源:Kafka(Receiver 与 Direct 方式)
Receiver 方式(已过时,不推荐)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| createStream | KafkaUtils.createStream(ssc, zkQuorum, groupID, topics) | 使用 Receiver 从 Kafka 消费数据 | val stream = KafkaUtils.createStream(ssc, "zk:2181", "group1", Map("topic" -> 1)) | 使用 ZooKeeper 管理 offset;Receiver 单点故障;不支持精确一次。 |
Direct 方式(推荐)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| createDirectStream | KafkaUtils.createDirectStream(ssc, kafkaParams, topics) | Direct 方式拉取 Kafka 数据,无 Receiver | val 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 可手动提交。 |
| 获取 offset | stream.foreachRDD { rdd => val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges // 处理逻辑 // 可在此提交 offset 到 Kafka 或外部存储} | 获取当前批次的 offset 范围 | 见上 | 需将 RDD 转换为 HasOffsetRanges 类型。 |
Direct 方式优势:
- 无单点故障
- 更好的容错性
- 支持从任意 offset 开始消费
- 可以与 checkpoint 结合实现精确一次处理
4.3 Flume 与 Kinesis 集成
Flume 集成
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| createStream | FlumeUtils.createStream(ssc, host, port) | 从 Flume agent 接收数据(Push-based) | val stream = FlumeUtils.createStream(ssc, "flume-agent", 41414) | 需 Flume agent 配置 Avro Source。 |
| createPollingStream | FlumeUtils.createPollingStream(ssc, hosts) | 轮询方式从多个 Flume agent 拉取数据 | val stream = FlumeUtils.createPollingStream(ssc, List("agent1", "agent2")) | 更可靠,支持负载均衡。 |
Kinesis 集成(AWS)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| createStream | KinesisUtils.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 凭证;适用于云环境数据采集。 |
| createDirectStream | KinesisUtils.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 常用窗口转换操作
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| window | dstream.window(windowLength, slideInterval) | 返回一个基于窗口的 DStream,后续可进行任意转换 | val windowedStream = stream.window(Seconds(30), Seconds(10)) | 返回新的 DStream,可链式调用 map、reduce 等操作。 |
| countByWindow | dstream.countByWindow(windowLength, slideInterval) | 统计窗口内元素总数 | val countStream = stream.countByWindow(Seconds(60), Seconds(10)) | 高效实现,避免手动 reduce。 |
| reduceByWindow | dstream.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)) | 每次重新计算整个窗口,性能较差。 |
| countByValueAndWindow | dstream.countByValueAndWindow(windowLength, slideInterval) | 按值统计窗口内各元素出现次数 | val counts = words.countByValueAndWindow(Seconds(60), Seconds(10)) | 类似 countByValue,但作用于窗口数据。 |
| union + window | stream1.union(stream2).window(...) | 合并多个流后进行窗口操作 | val combined = src1.union(src2).window(Seconds(20), Seconds(5)) | 可用于多源数据的统一窗口分析。 |
reduceByWindow 增量更新说明:
reduceFunc:正向聚合函数(如_ + _)invReduceFunc:反向减少函数(如_ - _),用于移除滑出窗口的数据- 若提供反向函数,Spark 可基于前一个窗口结果增量更新,极大提升性能
5.3 基于窗口的状态管理
| 概念/方法 | 说明 | 注意事项 |
|---|---|---|
| 窗口状态 | 窗口操作本身不保存跨窗口的状态,每次计算独立 | 若需跨多个窗口累积状态,需结合 updateStateByKey 或 mapWithState。 |
| 状态持久化 | 窗口计算结果可通过 foreachRDD 写入外部存储(如 Redis、数据库)实现持久化 | 避免在 driver 内存中累积大量状态,防止 OOM。 |
| 与 mapWithState 结合 | 可将窗口聚合结果作为输入,传递给 mapWithState 进行长期状态维护 | 例如:每 10 秒统计过去 1 分钟的 UV,再累计全天 PV。 |
| Checkpoint 支持 | 窗口操作本身不依赖 checkpoint,但若使用有状态操作(如 updateStateByKey),需启用 checkpoint | 必须设置 checkpoint 目录以支持故障恢复。 |
第六章:有状态流处理
6.1 updateStateByKey 操作
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| updateStateByKey | dstream.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 操作(推荐方式)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| mapWithState | dstream.mapWithState(stateSpec) | 高效的有状态映射操作,仅处理有数据输入的 key | val spec = StateSpec.function(updateFunc).timeout(Minutes(10))val mappedDStream = keyValueStream.mapWithState(stateSpec) | 推荐替代 updateStateByKey,性能更好。 |
| StateSpec.function | StateSpec.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)机制。 |
| timeout | spec.timeout(duration) | 设置状态空闲超时时间,超时后自动清除 | spec.timeout(Minutes(5)) | 减少内存占用,避免状态无限增长。 |
| Checkpoint | 同样需要启用 checkpoint | ssc.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)机制
| 方法/概念 | 语法/说明 | 用途 | 注意事项 |
|---|---|---|---|
| checkpoint | ssc.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 目录重建 StreamingContext | val 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 处理过多数据 | 可通过 repartition 或 coalesce 调整分区数。 |
8.3 数据序列化与内存管理
| 配置项 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| spark.serializer | spark.serializer org.apache.spark.serializer.KryoSerializer | 使用 Kryo 序列化,比 Java 序列化更高效 | 必须注册自定义类以获得最佳性能。 |
| spark.kryo.classesToRegister | spark.kryo.classesToRegister com.example.Event | 注册需要 Kryo 序列化的类 | 提高序列化速度,减少内存占用。 |
| spark.streaming.kafka.maxRatePerPartition | kafkaParams.put("maxRatePerPartition", "1000") | 限制每个 Kafka 分区每秒拉取的消息数 | 防止 executor 内存溢出(OOM),尤其在流量突增时。 |
| spark.streaming.backpressure.enabled | spark.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.bytes | kafkaParams.put("fetch.message.max.bytes", "10485880") | 单次拉取最大数据量(默认 1MB) | 可适当增大以提高吞吐,但需匹配 executor 内存。 |
| max.partition.fetch.bytes | kafkaParams.put("max.partition.fetch.bytes", "10485880") | 每个分区单次拉取最大字节数 | 建议与 fetch.message.max.bytes 一致。 |
| auto.offset.reset | kafkaParams.put("auto.offset.reset", "latest") | 无初始 offset 时从最新(latest)或最早(earliest)开始 | 生产环境建议设为 latest 避免重放历史数据。 |
| 手动提交 offset | 在 foreachRDD 中调用 commitAsync | 实现精确一次语义 | 需确保写入外部系统与 offset 提交在同一个语义下。 |
| 并行度匹配 | Kafka 分区数 ≥ Spark 分区数 | 确保最大并行消费能力 | 若 Kafka 分区少,Spark 无法并行处理。 |
第九章:与外部系统集成
9.1 写入 Kafka
| 方法/库 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
| 使用 foreachRDD(推荐) | 通过 foreachRDD 写入 Kafka,使用 KafkaProducer | dstream.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 建立连接,批量执行 SQL | dstream.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 API | Structured 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 + Alertmanager | CTR 异常下降时触发告警 |
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 Source | offsetLog 中的 currentOffset vs endOffset | 差值大表示消费滞后(lag) |
建议: 定期监控 UI,设置 Prometheus + Grafana 做长期趋势分析。