Quarkus结合SmallRye Kafka:如何检查Kafka存活并处理发送错误?
问题解答
一、检查Kafka服务器状态
Quarkus集成了MicroProfile Health规范,可通过添加健康检查扩展来监控Kafka连接状态:
- 添加健康检查依赖
<dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-smallrye-health</artifactId> </dependency>
- 启用Kafka健康检查
在application.properties中配置:
# 启用Kafka生产者健康检查 quarkus.smallrye-health.kafka.enabled=true # 可选:设置健康检查超时时间 quarkus.smallrye-health.kafka.timeout=5s
启动应用后,访问/q/health端点即可查看Kafka的健康状态,当Kafka不可用时,健康状态会标记为DOWN。
二、发送数据时的异常处理
Emitter.send()方法返回CompletionStage<Void>,可通过它处理发送失败的异常:
emitter.send(Record.of(hotel.getId(), hotel)) .whenComplete((unused, throwable) -> { if (throwable != null) { // 自定义失败处理逻辑,比如日志记录、重试触发、告警等 LOGGER.error("发送酒店数据到Kafka失败,ID: {}", hotel.getId(), throwable); } else { LOGGER.info("酒店数据发送成功,ID: {}", hotel.getId()); } });
也可通过配置实现自动重试,在application.properties中添加:
# 配置生产者重试次数 mp.messaging.outgoing.my-channel.kafka.producer.retries=3 # 重试间隔时间 mp.messaging.outgoing.my-channel.kafka.producer.retry.backoff.ms=1000
三、参考文档
- Quarkus官方文档中SmallRye Reactive Messaging Kafka章节,涵盖生产者配置、异常处理的详细说明
- Quarkus健康检查相关文档,介绍如何扩展和自定义健康检查逻辑
内容的提问来源于stack exchange,提问作者hunter662
相关产品推荐
相关产品推荐

