【发布时间】: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