Spark Streaming 原理
概述
Spark Streaming 是 Spark 早期的实时计算框架,核心思想是微批处理:把连续数据流切成小批次(Batch),交给 Spark 引擎执行。本文讲清 DStream 模型、Receiver 架构、批次调度机制,以及与 Storm/Flink 的定位差异。
一、核心模型:DStream
1.1 什么是 DStream
DStream(Discretized Stream,离散化流)是 Spark Streaming 的编程抽象:把连续流按时间切成一系列 RDD。
连续数据流
│ 按 Batch Interval 切分
▼
┌─────────────────────────────────────────────┐
│ DStream = 无限 RDD 序列 │
│ [RDD@t1] [RDD@t2] [RDD@t3] [RDD@t4] ... │
└─────────────────────────────────────────────┘| 特性 | 说明 |
|---|---|
| 时间切分 | 按固定间隔(如 2 秒)切批次 |
| 每批是 RDD | 复用 Spark 算子与容错 |
| 批间衔接 | 状态可跨批累积 |
1.2 DStream 的运算
DStream 算子作用于每个批次,本质是对每批 RDD 的 map:
scala
val lines = ssc.socketTextStream("host", 9999)
val words = lines.flatMap(_.split(" "))
val pairs = words.map(w => (w, 1))
val counts = pairs.reduceByKey(_ + _)
counts.print()| 算子 | 作用 |
|---|---|
| map/filter/flatMap | 逐 RDD 转换 |
| reduceByKey | 批内聚合 |
| updateStateByKey | 跨批状态更新 |
| window | 滑动窗口 |
二、Batch Interval 与批次调度
2.1 概念
| 概念 | 说明 |
|---|---|
| Batch Interval | 切批间隔(如 2s),决定延迟下限 |
| Job 生成 | 每批生成一个 Job |
| 调度 | Job 按 FIFO 进入 Spark 执行 |
streamingContext = new StreamingContext(sc, Seconds(2))2.2 批次耗时与积压
| 情况 | 结果 |
|---|---|
| 批次处理耗时 < 间隔 | 正常,有间隙 |
| 批次处理耗时 > 间隔 | 批次积压,延迟越来越大 |
| 监控指标 | 说明 |
|---|---|
| Processing Time | 批次实际处理耗时 |
| Scheduling Delay | 批次排队等待时间 |
| Total Delay | 总延迟 |
持续积压时需扩容或优化,否则延迟恶化。
三、Receiver 架构
3.1 数据接收方式
Spark 集群
┌─────────┐ 接收 ┌──────────────┐
│ 数据源 │ ──────►│ Receiver(在 │
│ Kafka/ │ │ Executor 中) │
│ Socket │ └──────┬───────┘
└─────────┘ │ 写入 BlockManager
▼
数据块(本 Executor 内存/磁盘)
│ 复制到其他节点(副本)
▼
批次 RDD 计算| 组件 | 职责 |
|---|---|
| Receiver | 在 Executor 中持续接收数据 |
| BlockManager | 存储数据块并复制 |
| 批次组装 | 定时把块组装成 RDD 输入 |
3.2 两种接收模式
| 模式 | 说明 | 语义 |
|---|---|---|
| 默认(At-Least-Once) | Receiver 接收后异步写 Kafka offset | 可能重复 |
| 手动管理 offset | 数据落盘后再提交 offset | 接近 Exactly-Once |
3.3 Receiver 故障
| 故障 | 恢复 |
|---|---|
| Receiver 重启 | 自动重启并恢复 |
| 数据丢失 | 开启 WAL(日志预写)后可从日志恢复 |
| Executor 崩溃 | Receiver 迁移到其他节点 |
四、窗口操作
4.1 窗口概念
window duration:窗口长度(如 10s)
slide duration:滑动步长(如 5s)
t1 t2 t3 t4 t5 t6
├── 窗口1(10s)──┤
├── 窗口2(10s)──┤scala
// 每 5 秒统计最近 10 秒的词频
val windowed = pairs
.reduceByKeyAndWindow(_ + _, Seconds(10), Seconds(5))4.2 窗口性能
| 优化 | 说明 |
|---|---|
| 增量聚合 | 复用前一窗口结果(加新减旧) |
| 窗口缓存 | 减少重复计算 |
| 大窗口 | 注意状态内存占用 |
五、与 Storm / Flink 对比
5.1 处理模型
| 框架 | 模型 | 延迟 | 吞吐 |
|---|---|---|---|
| Storm | 纯流式(逐条) | 毫秒级 | 低-中 |
| Spark Streaming | 微批(DStream) | 秒级 | 高 |
| Flink | 流式 + 批式统一 | 毫秒级 | 高 |
5.2 语义与容错
| 维度 | Storm | Spark Streaming | Flink |
|---|---|---|---|
| Exactly-Once | 难(需事务) | 支持(WAL+幂等) | 原生支持 |
| 状态管理 | 弱 | 支持(StateDStream) | 强(Keyed State) |
| 事件时间 | 不支持 | 有限 | 原生支持 |
| 背压 | 弱 | 有限(Receiver 模式) | 强 |
5.3 定位
| 阶段 | 结论 |
|---|---|
| 历史 | Spark Streaming 是 Spark 的实时补充 |
| 现状 | 新项目被 Structured Streaming 取代 |
| 对比 Flink | 延迟不如 Flink,事件时间支持有限 |
新实时项目优先选 Structured Streaming 或 Flink;Spark Streaming 适合熟悉 Spark、延迟要求秒级以上的场景。
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 批次持续积压 | 处理时间 > 间隔,扩容或优化 |
| 数据重复 | 默认至少一次语义,需幂等输出 |
| Receiver 数据丢失 | 开启 WAL(checkpoint 到可靠存储) |
| 延迟高 | 减小 Batch Interval 或换实时框架 |
| 窗口状态 OOM | 减小窗口或清理过期状态 |