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