You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何优雅检测Kafka连通性并设置布尔标记useKafka

我需要实现日志写入的降级逻辑:成功连通Kafka broker时将日志写入Kafka,若Kafka服务不可用则直接跳过该逻辑,不阻塞主业务流程。原本计划设置布尔变量useKafka标记Kafka可用状态,但最初的实现写法笨拙,还触发了编译错误。

最初实现代码

...
val props = new Properties()
props.put("bootstrap.servers", kafkaBroker)
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
props.put("request.timeout.ms", 200)

var useKafka: Boolean = true

// 此处写法非常不优雅
val producer = try {
  new KafkaProducer[String, String](props)
}
catch {
  useKafka = false
}
...

上述代码编译时抛出错误:value isDefinedAt is not a member of Unit。报错的核心原因是Scala的catch块要求传入偏函数,必须通过case语句匹配异常类型,直接在catch块内写赋值语句,会让编译器把catch块识别为返回Unit的普通代码块,不符合语法约定。

useKafka变量的后续使用逻辑如下:

if (useKafka) producer.send(new ProducerRecord[String, String](kafkaTopic, "cobol", logStr))
2022年6月2日更新

基于社区给出的方案调整了连通性检测逻辑:KafkaConsumer在连接失败时会直接抛出错误,比KafkaProducer更适合做连通性校验,搭配带超时的Future可以灵活控制检测等待时长,避免长时间阻塞主流程。修改后的代码如下:

val simpleConsumer = new org.apache.kafka.clients.consumer.KafkaConsumer[String, String](props)
val testKafka = Try {Await.ready(Future(simpleConsumer.listTopics), 200  milliseconds)}.toOption

useKafka = testKafka.nonEmpty

内容的提问来源于stack exchange,提问作者Lars Skaug

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.03 05:27:31