ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SparkStreaming 之 Direct 模式深度剖析

SparkStreaming 之 Direct 模式深度剖析 摘要前面讲接收数据时提过Receiver 模式是 Spark Streaming 早期的 Kafka 集成方式Direct 模式是它的替代者。这篇把 Direct 模式讲透它为什么没有 Receiver、offset 怎么自己管理、怎么靠offset 和 Job 结果绑定拿到 exactly-once以及 createDirectStream 的 LocationStrategies 和 ConsumerStrategies 两个参数怎么选。关键词Spark Streaming, Direct 模式, Kafka, offset, exactly-once, createDirectStream一、Direct 模式解决的是 Receiver 模式的老问题先回顾 Receiver 模式的两个硬伤offset 提交和处理分离。Receiver 用 Kafka 的 High-Level Consumer 拉数据offset 交给 ZooKeeper 自动提交。处理失败时 offset 可能已经提交了数据就丢了开了 WAL 变 at-least-once又可能重复消费。资源开销大。每个 Receiver 要占一个核还要开 WAL 写盘保证不丢。Direct 模式的思路很直接把 Receiver 拿掉。没有 Receiver就没有长驻进程 WAL ZK offset这一整套负担。二、没有 Receiver数据怎么进来Direct 模式每个 batch直接从 Kafka 分区拉取指定 offset 范围的数据拉完就算算完走人不留长驻进程。用代码看最直观valstreamKafkaUtils.createDirectStream[String,String](ssc,PreferConsistent,// 位置策略Subscribe[String,String](topics,kafkaParams)// 订阅策略)每个 batch 里Spark 会为每个 Kafka 分区确定一个 offset 区间fromOffset ~ untilOffset然后 Executor 直接去 Kafka 把这区间数据拉回来组成这个 batch 的 RDD。三、offset 自管理exactly-once 的由来这是 Direct 模式最核心的价值。offset 不再交给 ZooKeeper而是Spark 自己管理、存在 checkpoint 里并且offset 的更新和 Job 的成功与否绑定每个 RDD 记录它消费的[fromOffset, untilOffset]。只有这个 Job 真正处理成功了offset 才会被提交写进 checkpoint。Job 失败 → 不提交 offset → 重算时从原 offset 重新拉数据不丢也不重。这就是 exactly-once 的来源。Receiver 模式做不到这点因为它管不住 offset 的提交时机。四、分区对齐并行度由 Kafka 分区决定Direct 模式里RDD 的分区数 Kafka 的分区数每个 Kafka 分区对应一个 RDD 分区。这带来一个直接的推论想提高并行度就增加 Kafka 分区数反过来Executor 再多如果 Kafka 分区只有 3 个并行度也上不去一个 batch 最多 3 个 task。所以 Direct 模式下调并行度不能光看 executor 数得先看 topic 的分区数。这也是做容量规划时容易忽略的一点。五、两个参数LocationStrategies 和 ConsumerStrategiescreateDirectStream 的两个关键参数分别回答在哪算和算哪些。LocationStrategies位置策略决定 Kafka 分区数据放在哪个 Executor 上处理PreferConsistent分区均匀分布到所有 Executor。大多数场景用它。PreferBrokersExecutor 和 Kafka broker 同节点时用能走本地读。PreferFixed手动指定分区到主机的映射一般用不上。ConsumerStrategies订阅策略决定订阅哪些 topic/分区Subscribe订阅指定 topic 列表。最常用。SubscribePattern用正则匹配 topic。Assign显式指定要消费的具体分区。六、完整代码 三个坑importorg.apache.spark.streaming.kafka010._importorg.apache.kafka.common.serialization.StringDeserializervalkafkaParamsMap[String,Object](bootstrap.servers-broker1:9092,broker2:9092,key.deserializer-classOf[StringDeserializer],value.deserializer-classOf[StringDeserializer],group.id-streaming-app,auto.offset.reset-latest,enable.auto.commit-(false:java.lang.Boolean)// 必须 false)ssc.checkpoint(hdfs://namenode:8020/checkpoint/app)valstreamKafkaUtils.createDirectStream[String,String](ssc,PreferConsistent,Subscribe[String,String](Array(topic-a),kafkaParams))三个容易踩的坑enable.auto.commit必须设 false。Direct 模式的 offset 由 Spark 自己管如果不关掉 Kafka 的自动提交就会有两套 offset 在打架exactly-once 就废了。必须开 checkpoint。offset 存在 checkpoint 里不开 checkpointDriver 挂了 offset 就丢重启后不知道从哪继续。Kafka 分区数变了要重启。Direct 模式把 offset 按分区记在 checkpoint 里如果运行时新增/减少 Kafka 分区offset 对不上会报错需要重启应用并清理 checkpoint。七、总结Direct 模式去掉 Receiver直接按 offset 区间从 Kafka 分区拉数据省掉长驻进程、WAL、ZK 这套负担。offset 由 Spark 自管理、存 checkpoint且与 Job 成功绑定从而拿到 exactly-once。分区对齐意味着并行度 Kafka 分区数调并行度先看分区。LocationStrategies 常用 PreferConsistentConsumerStrategies 常用 Subscribe。三个坑auto.commit 设 false、必须开 checkpoint、分区数变更要重启。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
返回列表