如何优雅检测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
相关产品推荐
相关产品推荐

