Spark Streaming:从批处理走向流处理
0. 引言
Spark 凭借内存计算与 DAG 简化处理流程,取代 Hadoop MapReduce 成为批处理明星框架。流计算兴起后,Spark 将触角延伸至流计算领域,诞生了 Spark Streaming——建立在 Spark 批处理技术之上的流计算框架,提供可扩展、高吞吐、错误容忍的流式处理。本文从系统架构、流的描述、流的处理、流的状态、消息处理可靠性五个方面分析。
1. 系统架构:微批处理模型
Spark Streaming 本质上是将流数据分成一段段块数据后,进行连续不断的批处理——这就是"微批"(Micro-Batch)模型,与 Flink 的逐事件处理形成鲜明对比。
2. 流的描述:DStream 是 RDD 的模板
| 概念 | 对应抽象 | 说明 |
|---|---|---|
| RDD | 数据集合 | Spark 引擎的计算单元 |
| DStream | 队列(流) | 对流的抽象,内部是同一类 RDD 的模板,可视为 RDD 组成的序列 |
| Transformation | 处理节点 | map、flatMap、filter、reduce、union、join、transform、updateStateByKey 等 |
| DStreamGraph | DAG | 对流计算过程的描述,是 RDD 处理过程 DAG 的模板 |
| Output Operations | 输出 | print、saveAsTextFiles、saveAsHadoopFiles、foreachRDD;转化由输出操作触发执行 |
3. 流的处理
3.1 输入
三种输入流:基础数据源(socketTextStream、textFileStream、queueStream)、高级数据源(Kafka、Flume、Kinesis 等外部工具类)、自定义数据源(继承 Receiver 抽象类):
SparkConf conf = new SparkConf().setMaster("local[2]").setAppName("WordCountExample");
JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(1));
JavaReceiverInputDStream<String> lines = jssc.socketTextStream("localhost", 9999);3.2 处理与输出
三类转换操作:常用流式处理(map/filter/reduce/count/transform)、流数据状态操作(union/join/cogroup/window)、流信息状态操作(updateStateByKey/mapWithState):
JavaDStream<String> words = lines.flatMap(x -> Arrays.asList(x.split(" ")).iterator());
JavaPairDStream<String, Integer> pairs = words.mapToPair(s -> new Tuple2<>(s, 1));
JavaPairDStream<String, Integer> wordCounts = pairs.reduceByKey((i1, i2) -> i1 + i2);
wordCounts.print(); // 输出:print/saveAsTextFiles/saveAsHadoopFiles/foreachRDD3.3 反向压力:借鉴 PID 控制器
Spark 1.5 起引入反向压力(默认关闭,需设置 spark.streaming.backpressure.enabled=true),借鉴工业控制 PID 思路:
- 处理完每批数据后统计处理结束时间、时延、等待时延、消息数;
- 据此估计处理速度并通知数据生产者;
- 生产者动态调整生产速度,使生产与处理速度匹配。
4. 流的状态
- 流数据状态:组成 DStream 的 RDD 本身就是流数据状态;window API 实现窗口管理(count/reduce 聚合),union/join/cogroup 实现多流关联;
- 流信息状态:
updateStateByKey与mapWithState基于 key 记录历史信息并随新数据更新。区别:updateStateByKey 返回全部历史信息(完整直方图),mapWithState 只返回本批更新的信息(变化的柱条)——性能上mapWithState更优,实时场景更推荐。
5. 消息处理可靠性
| 环节 | 机制 | 保证级别 |
|---|---|---|
| 数据接收 | WAL(1.2+)写错误容忍存储 + 可靠接收器(Kafka) | at least once |
| 数据接收 | Kafka Direct API(1.3+) | exactly once |
| 数据处理 | 基于 RDD(天然幂等) | exactly once |
| 数据输出 | 无事务保障 | at least once(可能重复输出) |
输出重复的应对:为每批 RDD 增加唯一标识符,写入数据库时先检查是否已写入(配合事务保证完整性)。
6. 小结
- 模型:微批——流被切成 RDD 块连续批处理,最小间隔典型为 100 毫秒,延迟下限高、无法逐事件处理;
- 描述:DStream(RDD 模板)/ DStreamGraph(DAG 模板);
- 处理:三类转换操作 + PID 式反向压力;
- 状态:RDD 天然流数据状态;updateStateByKey(全量)vs mapWithState(增量,更优);
- 可靠性:接收端 WAL + Kafka Direct 可达 exactly once,输出端仅 at least once。
Spark Streaming 代表了"流批一体"趋势,但微批延迟与"不能逐条处理"的局限,使它对严苛实时场景并不适用。
版本注:Spark 2.3+ 引入 Structured Streaming(基于 DataFrame/SQL 的流处理,支持事件时间与端到端 exactly-once),DStream API 已进入维护模式;新项目建议使用 Structured Streaming。
下一章讲解 Apache Samza:最简洁的开源流计算框架。