【发布时间】:2020-07-22 12:25:25
【问题描述】:
我正在开发一个 Scala 应用程序。我在里面使用卡夫卡。我要使用来自 kafka 主题的消息。由于我正在编写一个测试用例,我需要记录做一些断言来通过我的测试用例。我正在使用以下代码来使用 kafka 消息:
import java.util.{Collections, Properties}
import java.util.regex.Pattern
import org.apache.kafka.clients.consumer.KafkaConsumer
import scala.collection.JavaConverters._
object KafkaConsumerSubscribeApp extends App {
val props:Properties = new Properties()
props.put("group.id", "test")
props.put("bootstrap.servers","localhost:9092")
props.put("key.deserializer",
"org.apache.kafka.common.serialization.StringDeserializer")
Props.put("value.deserializer",
"org.apache.kafka.common.serialization.StringDeserializer")
props.put("enable.auto.commit", "true")
props.put("auto.commit.interval.ms", "1000")
val consumer = new KafkaConsumer(props)
val topics = List("topic_text")
try {
consumer.subscribe(topics.asJava)
while (true) {
val records = consumer.poll(10)
for (record <- records.asScala) {
println("Topic: " + record.topic() +
",Key: " + record.key() +
",Value: " + record.value() +
", Offset: " + record.offset() +
", Partition: " + record.partition())
}
}
}catch{
case e:Exception => e.printStackTrace()
}finally {
consumer.close()
}
}
这段代码我面临两个问题。在 intellij 中,它警告 poll 方法已被弃用。如何修改此代码以弃用?第二个问题是我希望这个方法返回它从 kafka 主题获得的消息。该消息在record.value() 中。我该如何退货?使用此代码,因为使用了 while(true),所以它将是一个无限循环,它将继续侦听来自主题的消息。如何从该方法返回 record.value() 以便可以在其他方法中使用从主题获取的数据。
【问题讨论】:
标签: scala apache-kafka