【问题标题】:Is it possible to obtain specific message offset in Kafka+SparkStreaming?是否可以在 Kafka+Spark Streaming 中获取特定的消息偏移量?
【发布时间】:2016-10-01 01:19:25
【问题描述】:

我正在尝试使用 Spark Direct Stream 在 Kafka 中获取和存储特定消息的偏移量。 查看 Spark 文档很容易获得每个分区的范围偏移量,但我需要的是在完全扫描队列后存储主题的每条消息的起始偏移量。

【问题讨论】:

    标签: apache-spark apache-kafka spark-streaming kafka-consumer-api


    【解决方案1】:

    是的,您可以使用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”是扫描分区的起始偏移量。非常感谢avr
    • 是的,你的假设是正确的。 fromOffset 是消费者开始阅读的起始偏移量。
    • 你如何管理偏移量,所以我将 inputOffset 设置为 0,但我让我的 spark 流应用程序运行,然后它崩溃了,它会再次从 0 开始正确吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-03-13
    • 2017-06-22
    • 2017-02-06
    • 2018-09-22
    • 1970-01-01
    • 1970-01-01
    • 2021-05-22
    相关资源
    最近更新 更多