Spark Streaming 与 Kafka 集成
概述
Kafka 是流处理事实标准的数据管道,Spark Streaming 消费 Kafka 有 Receiver 与 Direct 两种方式。本文以 Direct 方式为主线:offset 管理、背压机制、端到端语义保证,以及生产调优参数。
一、两种消费方式
1.1 Receiver 方式(旧)
Kafka → Receiver(Executor)→ WAL → BlockManager → 批次 RDD
offset 存在 ZooKeeper,与数据处理不同步| 特点 | 说明 |
|---|---|
| 数据先收后算 | Receiver 在 Executor 内接收 |
| offset 异步提交 | 崩溃可能重复/丢失 |
| 需 WAL | 开启 WAL 保数据 |
| 语义 | At-Least-Once |
1.2 Direct 方式(推荐)
每个批次:Driver 直接向 Kafka 拉取对应 offset 范围的数据
Kafka → 批次 RDD(按 offset 范围)| 特点 | 说明 |
|---|---|
| 分区对齐 | RDD 分区 = Kafka 分区 |
| 无中间层 | 直接拉取,无需 WAL |
| offset 自主管理 | 存在 checkpoint |
| 语义 | 可精确一次 |
| 对比 | Receiver | Direct |
|---|---|---|
| 架构 | Receiver 中转 | 直连 Kafka |
| 并行度 | 与 Receiver 数相关 | 与分区数一致 |
| 容错 | 依赖 WAL | 按 offset 重放 |
| 推荐 | 不推荐 | 生产标准 |
二、Direct 方式原理
2.1 流程
1. 每个批次:Driver 获取各分区最新 offset
2. 确定本批拉取范围 [startOffset, endOffset]
3. 生成 RDD,分区与 Kafka 分区一一对应
4. Executor 按范围拉取消费
5. 处理完成,offset 存入 checkpoint2.2 代码
scala
val stream = KafkaUtils.createDirectStream[String, String](
ssc,
LocationStrategies.PreferConsistent,
ConsumerStrategies.Subscribe[String, String](
List("orders"),
Map("bootstrap.servers" -> "kafka:9092")
)
)| 参数 | 说明 |
|---|---|
maxRatePerPartition | 每分区每秒消费上限(背压) |
enable.auto.commit | 设为 false,手动管理 offset |
group.id | 消费组 |
三、offset 管理
3.1 为什么重要
offset 是"消费到哪里"的指针,管理不当会导致丢数据或重复。
3.2 管理方式
| 方式 | 说明 |
|---|---|
| checkpoint | 引擎自动存 offset(推荐) |
| Kafka 自身 | 提交到 __consumer_offsets |
| ZooKeeper | 旧方式 |
| 外部存储 | MySQL/Redis 自主管理 |
3.3 checkpoint 恢复
故障重启:
从 checkpoint 恢复 offset
→ 从上次位置继续消费
→ 未处理数据不丢| 注意 | 说明 |
|---|---|
| 幂等输出 | checkpoint 恢复会重放部分数据 |
| checkpoint 目录 | 每应用独立,勿共用 |
四、背压机制
4.1 什么是背压
消费速度跟不上生产速度时,反向限制消费速率,防止积压失控与 OOM。
4.2 Spark Streaming 背压
| 参数 | 默认 | 说明 |
|---|---|---|
spark.streaming.backpressure.enabled | false | 开启背压 |
spark.streaming.backpressure.initialRate | - | 初始速率 |
spark.streaming.kafka.maxRatePerPartition | 0(无限) | 每分区速率上限 |
开启后:
根据批次处理耗时动态调整消费速率
处理快 → 加速;处理慢 → 减速4.3 与积压的关系
| 情况 | 处理 |
|---|---|
| 背压生效 | 速率自动收敛到可处理水平 |
| 仍积压 | 扩容 Executor 或优化作业 |
| 峰值数据 | 背压平滑,避免雪崩 |
五、端到端语义保证
5.1 语义链条
Kafka 源(可重放)
→ Direct 消费(offset 精确)
→ 计算(checkpoint 恢复)
→ 输出(幂等/事务)
= 端到端 Exactly-Once5.2 各端策略
| 端 | 保证手段 |
|---|---|
| 源 | Kafka 消息保留期内可重放 |
| 消费 | checkpoint 记录 offset |
| 计算 | 状态 + 结果 checkpoint |
| 输出 | 幂等写或事务提交 |
5.3 常见组合
| 组合 | 语义 |
|---|---|
| Direct + checkpoint + 文件 Sink | 精确一次 |
| Direct + checkpoint + 唯一键写 MySQL | 等效精确一次 |
| 无幂等输出 | 至少一次 |
六、生产调优
6.1 并行度
| 要点 | 说明 |
|---|---|
| 分区对齐 | 增加 Kafka 分区可提升并行 |
| RDD 分区 | Direct 下等于 Kafka 分区 |
| 重分区 | 需要更多并行可 repartition |
6.2 参数汇总
| 参数 | 说明 |
|---|---|
spark.streaming.kafka.maxRatePerPartition | 每分区速率上限 |
spark.streaming.backpressure.enabled | 背压开关 |
spark.streaming.backpressure.rateEstimator | 速率估算器 |
spark.streaming.blockInterval | 块切分间隔 |
spark.streaming.receiver.writeAheadLog.enable | WAL(Receiver 用) |
6.3 监控指标
| 指标 | 说明 |
|---|---|
| Input Rate | 输入速率 |
| Processing Time | 批次处理耗时 |
| 积压量 | Kafka 未消费消息数 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 消费不均衡 | 分区数变化,检查分区对齐 |
| offset 重置 | checkpoint 丢失,用 auto.offset.reset 控制 |
| 重复消费 | 输出未幂等,或 checkpoint 与幂等缺失 |
| 积压暴涨 | 开启背压,或扩容优化 |
| Direct 比 Receiver 慢 | 分区数少,增加 Kafka 分区 |