Skip to content

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 这类二进制兼容问题。

scala
// build.sbt
libraryDependencies += "org.apache.spark" %% "spark-streaming-kafka-0-10" % "2.4.0"
xml
<!-- pom.xml -->
<dependency>
  <groupId>org.apache.spark</groupId>
  <artifactId>spark-streaming-kafka-0-10_2.11</artifactId>
  <version>2.4.0</version>
</dependency>

二、KafkaParams 配置

Direct 模式通过一个 Map[String, Object] 传 Kafka 消费参数。几个必填的:

scala
import org.apache.kafka.common.serialization.StringDeserializer

val kafkaParams = Map[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 调用

架构图

scala
import org.apache.spark.streaming.kafka010._

val stream = KafkaUtils.createDirectStream[String, String](
  ssc,
  PreferConsistent,                                    // 位置策略
  Subscribe[String, String](Array("topic-a"), kafkaParams)  // 订阅策略
)

两个参数分别是位置策略(数据放哪算)和订阅策略(订阅哪些分区)。


四、三种 ConsumerStrategies

架构图

订阅策略决定"消费哪些 topic/分区":

scala
// 1. Subscribe:订阅指定 topic 列表(最常用)
Subscribe[String, String](Array("topic-a", "topic-b"), kafkaParams)

// 2. SubscribePattern:正则匹配 topic
SubscribePattern[String, String](
  java.util.regex.Pattern.compile("topic-.*"), kafkaParams)

// 3. Assign:显式指定具体分区
Assign[String, String](
  Array(new TopicPartition("topic-a", 0), new TopicPartition("topic-a", 1)),
  kafkaParams)

Subscribegroup.id 自动分配分区,是最常见的用法;Assign 手动指定分区,用在需要精确控制消费范围的场景(比如回放某个分区的一段 offset)。


五、三种 LocationStrategies

位置策略决定"分区数据放到哪个 Executor 处理":

scala
// 1. PreferConsistent:分区均匀分布到所有 Executor(最常用)
LocationStrategies.PreferConsistent

// 2. PreferBrokers:Executor 和 Kafka broker 同节点时用
LocationStrategies.PreferBrokers

// 3. PreferFixed:手动指定分区到主机的映射
LocationStrategies.PreferFixed(
  Map(new TopicPartition("topic-a", 0) -> "host1:9092"))

绝大多数场景用 PreferConsistent。只有当你把 Executor 部署在 Kafka broker 同一台机器上时,PreferBrokers 才有意义(能走本地读)。PreferFixed 一般用不上。


六、offset 手动管理

这是 Direct 模式代码里最容易写错的部分。要手动提交 offset,得先拿到每个 RDD 对应的 offset 范围:

scala
import org.apache.spark.streaming.kafka010._

stream.foreachRDD { rdd =>
  // 关键:只能在 Direct stream 的第一个转换里拿 offsetRanges
  val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges

  // 处理 rdd ...
  val result = rdd.map(...).reduceByKey(...)

  // 处理成功后再提交 offset
  stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)
}

两个坑:

  1. offsetRanges 只能在第一个转换里取HasOffsetRanges 这个类型只在 Direct stream 生成的原始 RDD 上存在,一旦你 rdd.map(...) 之后,类型就丢了,再取会抛 ClassCastException。所以要在 foreachRDD 一进来就把 offsetRanges 抓出来。
  2. 处理成功再提交commitAsync 要放在处理逻辑之后,这样才能实现"处理—提交"的原子性——处理失败就不会提交 offset,重算时从原位置重拉,这正是 exactly-once 的实现方式。

七、完整代码串起来

scala
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.streaming.kafka010._
import org.apache.kafka.common.serialization.StringDeserializer

def createContext(): StreamingContext = {
  val ssc = new StreamingContext(conf, Seconds(5))
  ssc.checkpoint("hdfs://namenode:8020/checkpoint/app")

  val kafkaParams = Map[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)
  )

  val stream = KafkaUtils.createDirectStream[String, String](
    ssc, PreferConsistent,
    Subscribe[String, String](Array("topic-a"), kafkaParams))

  stream.foreachRDD { rdd =>
    val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
    rdd.map(_.value).foreachPartition { iter => /* 业务处理 */ }
    stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)
  }

  ssc
}

val ssc = StreamingContext.getOrCreate(checkpointPath, createContext _)
ssc.start(); ssc.awaitTermination()

提交时用 cluster 模式 + supervise(配合上一篇 Driver HA 的内容):

bash
spark-submit --master yarn --deploy-mode cluster --supervise \
  --class com.example.DirectStreamApp app.jar

八、总结

  • 依赖用 spark-streaming-kafka-0-10,版本与 Spark 严格对齐。
  • KafkaParams 里 enable.auto.commit 必须设 false,且 Scala 里要写 (false: java.lang.Boolean)
  • ConsumerStrategies 常用 Subscribe,LocationStrategies 常用 PreferConsistent。
  • offset 手动管理的关键:offsetRanges 只能在第一个转换里取,处理成功后再 commitAsync
  • 提交侧 cluster 模式 + supervise,配合 checkpoint 实现完整的高可用。

作者:大数据技术实践者
博客blog.starzy.cn
GitHubstarzy1990.github.io
专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践