【问题标题】:If Kafka Consumer fails (Spark Job), how to fetch the last offset committed by Kafka Consumer. (Scala)如果 Kafka Consumer 失败(Spark Job),如何获取 Kafka Consumer 提交的最后一个偏移量。 (斯卡拉)
【发布时间】:2018-08-14 15:33:59
【问题描述】:

在我提供任何细节之前,请注意,我不是询问如何使用 kafka-run-class.sh kafka.tools.ConsumerOffsetChecker 从控制台获取最新的偏移量。

我正在尝试使用 Scala(2.11.8)在 Spark(2.3.1)中创建一个 kafka 消费者(kafka 版本 0.10),这将是容错的。通过容错,我的意思是,如果由于某种原因 kafka 消费者死亡并重新启动,它应该从最后一个偏移量恢复消费消息。

为了实现这一点,我使用以下代码提交 Kafka 偏移量,

    val kafkaParams = Map[String, Object](
"bootstrap.servers" -> "localhost:9092",
"key.deserializer" -> classOf[StringDeserializer],
"value.deserializer" -> classOf[StringDeserializer],
"group.id" -> "group_101",
"auto.offset.reset" -> "latest",
"enable.auto.commit" -> (false: java.lang.Boolean), /*because messages successfully polled by the consumer may not yet have resulted in a Spark output operation*/
"session.timeout.ms" -> (30000: java.lang.Integer),
"heartbeat.interval.ms" -> (3000: java.lang.Integer)
)

val topic = Array("topic_1")

val offsets = Map(new org.apache.kafka.common.TopicPartition("kafka_cdc_1", 0) -> 2L) /*Edit: Added code to fetch offset*/

val kstream = KafkaUtils.createDirectStream[String, String](
ssc,
PreferConsistent,
Subscribe[String, String](topic, kafkaParams, offsets)  /*Edit: Added offset*/ 
)

kstream.foreachRDD{ rdd =>
val offsetRange = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
if(!rdd.isEmpty()) {
  val rawRdd = rdd.map(record => 
 (record.key(),record.value())).map(_._2).toDS()
  val df = spark.read.schema(tabSchema).json(rawRdd)
  df.createOrReplaceTempView("temp_tab")
  df.write.insertInto("hive_table")
}
kstream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRange) /*Doing Async Commit Here */
}

我已经尝试了很多方法来获取给定主题的最新偏移量,但无法正常工作。

谁能帮我用 scala 代码来实现这一点,好吗?

编辑: 在上面的代码中,我试图通过使用来获取最后一个偏移量

val offsets = Map(new org.apache.kafka.common.TopicPartition("kafka_cdc_1", 0) -> 2L) /*Edit: Added code to fetch offset*/

但是上面代码获取的偏移量是0,不是最新的。无论如何要获取最新的偏移量吗?

【问题讨论】:

  • 我看到你有足够的代码来处理这个问题。您能发布您看到的错误或问题吗?
  • @AbhishekN 我已经编辑了问题以添加我用来获取最后一个偏移量的代码。但不知何故,我总是以使用此方法获取偏移量 0 结束。我想获取最后一个偏移量。你知道有什么方法可以实现吗?

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


【解决方案1】:

找到上述问题的解决方案。这里是。希望对有需要的人有所帮助。

语言:Scala、Spark Job

val kafkaParams = Map[String, Object](
"bootstrap.servers" -> "localhost:9092",
"key.deserializer" -> classOf[StringDeserializer],
"value.deserializer" -> classOf[StringDeserializer],
"group.id" -> "group_101",
"auto.offset.reset" -> "latest",
"enable.auto.commit" -> (false: java.lang.Boolean), /*because messages successfully polled by the consumer may not yet have resulted in a Spark output operation*/
"session.timeout.ms" -> (30000: java.lang.Integer),
"heartbeat.interval.ms" -> (3000: java.lang.Integer)
)

import java.util.Properties

//create a new properties object with Kafaka Parameters as done previously. Note: Both needs to be present. We will use the proprty object just to fetch the last offset

val kafka_props = new Properties()
kafka_props.put("bootstrap.servers", "localhost:9092")
kafka_props.put("key.deserializer",classOf[StringDeserializer])
kafka_props.put("value.deserializer",classOf[StringDeserializer])
kafka_props.put("group.id","group_101")
kafka_props.put("auto.offset.reset","latest")
kafka_props.put("enable.auto.commit",(false: java.lang.Boolean))
kafka_props.put("session.timeout.ms",(30000: java.lang.Integer))
kafka_props.put("heartbeat.interval.ms",(3000: java.lang.Integer))

val topic = Array("topic_1")

/*val offsets = Map(new org.apache.kafka.common.TopicPartition("topic_1", 0) -> 2L) Edit: Added code to fetch offset*/

val topicAndPartition = new org.apache.kafka.common.TopicPartition("topic_1", 0) //Using 0 as the partition because this topic does not have any partitions
val consumer = new KafkaConsumer[String,String](kafka_props)    //create a 2nd consumer to fetch last offset
import java.util
consumer.subscribe(util.Arrays.asList("topic_1"))   //Subscribe to the 2nd consumer. Without this step, the offsetAndMetadata can't be fetched.
val offsetAndMetadata = consumer.committed(topicAndPartition)    //Find last committed offset for the given topicAndPartition
val endOffset = offsetAndMetadata.offset().toLong   //fetch the last committed offset from offsetAndMetadata and cast it to Long data type.

val fetch_from_offset = Map(new org.apache.kafka.common.TopicPartition("topic_1", 0) -> endOffset) // create a Map with data type (TopicPartition, Long)

val kstream = KafkaUtils.createDirectStream[String, String](
ssc,
PreferConsistent,
Subscribe[String, String](topic, kafkaParams, fetch_from_offset) //Pass the offset Map of datatype (TopicPartition, Long) created eariler
)

kstream.foreachRDD{ rdd =>
val offsetRange = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
if(!rdd.isEmpty()) {
  val rawRdd = rdd.map(record => 
 (record.key(),record.value())).map(_._2).toDS()
  val df = spark.read.schema(tabSchema).json(rawRdd)
  df.createOrReplaceTempView("temp_tab")
  df.write.insertInto("hive_table")
}
kstream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRange) /*Doing Async offset Commit Here */
}

【讨论】:

    猜你喜欢
    • 2018-07-01
    • 1970-01-01
    • 2020-12-28
    • 2023-03-02
    • 1970-01-01
    • 2021-01-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多