【问题标题】:Kafka consumer consuming messages multiple times when trying to process the messages using SparkKafka 消费者在尝试使用 Spark 处理消息时多次消费消息
【发布时间】:2020-11-11 10:48:54
【问题描述】:

我有一个 Kafka 消费者,它从主题中读取消息并使用 spark 将其写入配置单元表。当我在 Yarn 上运行代码时,它会多次读取相同的消息。我在该主题中有大约 100,000 条消息。但是,我的消费者继续多次阅读相同的内容。当我做一个不同的时候,我得到了实际的计数。

这是我编写的代码。我想知道我是否缺少任何设置。

 val spark = SparkSession.builder()
      .appName("Kafka Consumer")
      .enableHiveSupport()
      .getOrCreate()

    import spark.implicits._

    val kafkaConsumerProperty = new Properties()
    kafkaConsumerProperty.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "---")
    kafkaConsumerProperty.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer")
    kafkaConsumerProperty.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer")
    kafkaConsumerProperty.put(ConsumerConfig.GROUP_ID_CONFIG, "draw_attributes")
    kafkaConsumerProperty.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest")
    kafkaConsumerProperty.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true")
    val topic = "space_orchestrator"
    val kafkaConsumer = new KafkaConsumer[String,String](kafkaConsumerProperty)
    kafkaConsumer.subscribe(Collections.singletonList(topic))
    
    while(true){

      val recordSeq = kafkaConsumer.poll(10000).toSeq.map( x => x.value())
      if(!recordSeq.isEmpty)
        {
          val newDf = spark.read.json(recordSeq.toDS)
          newDf.write.mode(SaveMode.Overwrite).saveAsTable("dmart_dev.draw_attributes")
        }
    }

【问题讨论】:

  • 出于兴趣,您的消费者阅读这 100,000 条消息的速度有多快?如果它小于自动提交间隔(默认为 5 秒,IIRC)并且您在 5 秒过去之前停止轮询而没有明确关闭消费者,则偏移量将永远不会提交(请参阅stackoverflow.com/questions/38230862/… 的答案)
  • 这种行为是我总是手动提交偏移量的众多原因之一。
  • @LeviRamsey:我试过手动提交。但是,我仍然看到重复的记录。
  • 当您说重复时,您的意思是两次读取相同的偏移量还是在多个偏移量处读取相同的消息。一般来说,Kafka 不能防止后者(它提供的保护通常会严重损害性能而无法真正有用)。
  • 我认为它两次读取相同的偏移量。

标签: scala apache-spark hadoop apache-kafka kafka-consumer-api


【解决方案1】:

或者,尝试手动设置偏移量。为此,应禁用自动提交 (enable.auto.commit = false)。对于手动提交,KafkaConsumers 提供了两种方法,即commitSync()commitAsync()。顾名思义,commitSync() 是一个阻塞调用,在偏移成功提交后返回,而 commitAsync() 立即返回。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-09-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-10-19
    • 2020-08-11
    • 2020-12-18
    • 2021-10-11
    相关资源
    最近更新 更多