{T}

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 等
DStreamGraphDAG对流计算过程的描述,是 RDD 处理过程 DAG 的模板
Output Operations输出print、saveAsTextFiles、saveAsHadoopFiles、foreachRDD;转化由输出操作触发执行

3. 流的处理

3.1 输入

三种输入流:基础数据源(socketTextStreamtextFileStreamqueueStream)、高级数据源(Kafka、Flume、Kinesis 等外部工具类)、自定义数据源(继承 Receiver 抽象类):

java
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):

java
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/foreachRDD

3.3 反向压力:借鉴 PID 控制器

Spark 1.5 起引入反向压力(默认关闭,需设置 spark.streaming.backpressure.enabled=true),借鉴工业控制 PID 思路:

  1. 处理完每批数据后统计处理结束时间、时延、等待时延、消息数;
  2. 据此估计处理速度并通知数据生产者;
  3. 生产者动态调整生产速度,使生产与处理速度匹配。

4. 流的状态

  • 流数据状态:组成 DStream 的 RDD 本身就是流数据状态;window API 实现窗口管理(count/reduce 聚合),union/join/cogroup 实现多流关联;
  • 流信息状态updateStateByKeymapWithState 基于 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:最简洁的开源流计算框架。