【发布时间】:2016-10-01 01:19:25
【问题描述】:
我正在尝试使用 Spark Direct Stream 在 Kafka 中获取和存储特定消息的偏移量。 查看 Spark 文档很容易获得每个分区的范围偏移量,但我需要的是在完全扫描队列后存储主题的每条消息的起始偏移量。
【问题讨论】:
标签: apache-spark apache-kafka spark-streaming kafka-consumer-api
我正在尝试使用 Spark Direct Stream 在 Kafka 中获取和存储特定消息的偏移量。 查看 Spark 文档很容易获得每个分区的范围偏移量,但我需要的是在完全扫描队列后存储主题的每条消息的起始偏移量。
【问题讨论】:
标签: apache-spark apache-kafka spark-streaming kafka-consumer-api
是的,您可以使用createDirectStream 的MessageAndMetadata 版本,它允许您访问message metadata。
您可以在此处找到返回 tuple3 的 Dstream 的示例。
val ssc = new StreamingContext(sparkConf, Seconds(10))
val kafkaParams = Map[String, String]("metadata.broker.list" -> (kafkaBroker))
var fromOffsets = Map[TopicAndPartition, Long]()
val topicAndPartition: TopicAndPartition = new TopicAndPartition(kafkaTopic.trim, 0)
val topicAndPartition1: TopicAndPartition = new TopicAndPartition(kafkaTopic1.trim, 0)
fromOffsets += (topicAndPartition -> inputOffset)
fromOffsets += (topicAndPartition1 -> inputOffset1)
val messagesDStream = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder, Tuple3[String, Long, String]](ssc, kafkaParams, fromOffsets, (mmd: MessageAndMetadata[String, String]) => {
(mmd.topic ,mmd.offset, mmd.message().toString)
})
在上述示例中,tuple3._1 将具有 topic,tuple3._2 将具有 offset,tuple3._3 将具有 message。
希望这会有所帮助!
【讨论】:
messagesDStream 中的每条消息相关联的偏移量。我的意思是createDirectStream 给你Dstream 的Tuple3 并且在每个元组中你会得到topic-name 和message 及其关联的offset。
fromOffset 是消费者开始阅读的起始偏移量。