ARTICLE DETAIL

资讯详情

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

SparkStreaming 之 Direct 模式 API 代码实现

SparkStreaming 之 Direct 模式 API 代码实现 摘要上一篇把 Direct 模式的原理讲清楚了这篇落到代码——从 Maven 依赖、KafkaParams 配置、createDirectStream 调用到三种 ConsumerStrategies 和 LocationStrategies 的写法再到 offset 的手动管理给出一套能直接抄来跑的生产代码。关键词Spark Streaming, Direct 模式, createDirectStream, KafkaParams, offset 管理, commitAsync一、先对齐依赖版本Direct 模式Kafka 0.10 API用的依赖是spark-streaming-kafka-0-10不是老的-0-8。版本必须和你的 Spark 版本严格对齐否则运行时抛NoSuchMethodError这类二进制兼容问题。// build.sbtlibraryDependenciesorg.apache.spark%%spark-streaming-kafka-0-10%2.4.0!-- pom.xml --dependencygroupIdorg.apache.spark/groupIdartifactIdspark-streaming-kafka-0-10_2.11/artifactIdversion2.4.0/version/dependency二、KafkaParams 配置Direct 模式通过一个Map[String, Object]传 Kafka 消费参数。几个必填的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)enable.auto.commit必须设 false。上一篇讲过Direct 模式的 offset 由 Spark 自己管如果这里不关掉 Kafka 的自动提交会有两套 offset 在打架exactly-once 就失效了。注意 Scala 里要显式写成(false: java.lang.Boolean)否则会被当成Boolean装箱类型对不上。三、createDirectStream 调用importorg.apache.spark.streaming.kafka010._valstreamKafkaUtils.createDirectStream[String,String](ssc,PreferConsistent,// 位置策略Subscribe[String,String](Array(topic-a),kafkaParams)// 订阅策略)两个参数分别是位置策略数据放哪算和订阅策略订阅哪些分区。四、三种 ConsumerStrategies订阅策略决定消费哪些 topic/分区// 1. Subscribe订阅指定 topic 列表最常用Subscribe[String,String](Array(topic-a,topic-b),kafkaParams)// 2. SubscribePattern正则匹配 topicSubscribePattern[String,String](java.util.regex.Pattern.compile(topic-.*),kafkaParams)// 3. Assign显式指定具体分区Assign[String,String](Array(newTopicPartition(topic-a,0),newTopicPartition(topic-a,1)),kafkaParams)Subscribe由group.id自动分配分区是最常见的用法Assign手动指定分区用在需要精确控制消费范围的场景比如回放某个分区的一段 offset。五、三种 LocationStrategies位置策略决定分区数据放到哪个 Executor 处理// 1. PreferConsistent分区均匀分布到所有 Executor最常用LocationStrategies.PreferConsistent// 2. PreferBrokersExecutor 和 Kafka broker 同节点时用LocationStrategies.PreferBrokers// 3. PreferFixed手动指定分区到主机的映射LocationStrategies.PreferFixed(Map(newTopicPartition(topic-a,0)-host1:9092))绝大多数场景用PreferConsistent。只有当你把 Executor 部署在 Kafka broker 同一台机器上时PreferBrokers才有意义能走本地读。PreferFixed一般用不上。六、offset 手动管理这是 Direct 模式代码里最容易写错的部分。要手动提交 offset得先拿到每个 RDD 对应的 offset 范围importorg.apache.spark.streaming.kafka010._ stream.foreachRDD{rdd// 关键只能在 Direct stream 的第一个转换里拿 offsetRangesvaloffsetRangesrdd.asInstanceOf[HasOffsetRanges].offsetRanges// 处理 rdd ...valresultrdd.map(...).reduceByKey(...)// 处理成功后再提交 offsetstream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)}两个坑offsetRanges只能在第一个转换里取。HasOffsetRanges这个类型只在 Direct stream 生成的原始 RDD 上存在一旦你rdd.map(...)之后类型就丢了再取会抛ClassCastException。所以要在 foreachRDD 一进来就把 offsetRanges 抓出来。处理成功再提交。commitAsync要放在处理逻辑之后这样才能实现处理—提交的原子性——处理失败就不会提交 offset重算时从原位置重拉这正是 exactly-once 的实现方式。七、完整代码串起来importorg.apache.spark.streaming.{Seconds,StreamingContext}importorg.apache.spark.streaming.kafka010._importorg.apache.kafka.common.serialization.StringDeserializerdefcreateContext():StreamingContext{valsscnewStreamingContext(conf,Seconds(5))ssc.checkpoint(hdfs://namenode:8020/checkpoint/app)valkafkaParamsMap[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))valstreamKafkaUtils.createDirectStream[String,String](ssc,PreferConsistent,Subscribe[String,String](Array(topic-a),kafkaParams))stream.foreachRDD{rddvaloffsetRangesrdd.asInstanceOf[HasOffsetRanges].offsetRanges rdd.map(_.value).foreachPartition{iter/* 业务处理 */}stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)}ssc}valsscStreamingContext.getOrCreate(checkpointPath,createContext _)ssc.start();ssc.awaitTermination()提交时用 cluster 模式 supervise配合上一篇 Driver HA 的内容spark-submit--masteryarn--deploy-mode cluster--supervise\--classcom.example.DirectStreamApp app.jar八、总结依赖用spark-streaming-kafka-0-10版本与 Spark 严格对齐。KafkaParams 里enable.auto.commit必须设 false且 Scala 里要写(false: java.lang.Boolean)。ConsumerStrategies 常用 SubscribeLocationStrategies 常用 PreferConsistent。offset 手动管理的关键offsetRanges只能在第一个转换里取处理成功后再commitAsync。提交侧 cluster 模式 supervise配合 checkpoint 实现完整的高可用。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
返回列表