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

如何通过application.yml配置为Camel ProducerTemplate添加Kafka健康检查?

为Camel Kafka ProducerTemplate添加健康检查的解决方案

你配置的camel.health.enabled: true和producersEnabled: true是开启Camel生产者健康检查的基础配置,但它们仅对Camel路由中显式定义的生产者端点生效——手动实例化的ProducerTemplate不在默认的健康检查覆盖范围内,这就是你看不到相关检查项的原因。

下面提供两种可行的解决方式:

方式一:自定义健康检查指示器

手动编写Spring Boot的HealthIndicator,验证ProducerTemplate与Kafka集群的连通性:

  1. 创建自定义健康检查类:
import org.apache.camel.ProducerTemplate;
import org.springframework.boot.actuate.health.Health;
import org.springframework.boot.actuate.health.HealthIndicator;
import org.springframework.kafka.core.KafkaAdmin;
import org.springframework.stereotype.Component;

@Component
public class KafkaProducerHealthIndicator implements HealthIndicator {

    private final ProducerTemplate producerTemplate;
    private final KafkaAdmin kafkaAdmin;

    // 注入所需依赖
    public KafkaProducerHealthIndicator(ProducerTemplate producerTemplate, KafkaAdmin kafkaAdmin) {
        this.producerTemplate = producerTemplate;
        this.kafkaAdmin = kafkaAdmin;
    }

    @Override
    public Health health() {
        try {
            // 用KafkaAdmin验证集群连通性
            kafkaAdmin.describeCluster();
            
            // 可选:用ProducerTemplate发送测试消息(需确保存在对应测试topic)
            // producerTemplate.sendBody("kafka:health-check-topic", "health-check-message");
            
            return Health.up()
                    .withDetail("kafka-producer-status", "连通正常")
                    .withDetail("cluster-id", kafkaAdmin.describeCluster().clusterId())
                    .build();
        } catch (Exception e) {
            return Health.down()
                    .withDetail("kafka-producer-status", "连通失败")
                    .withDetail("error-message", e.getMessage())
                    .build();
        }
    }
}
  1. 重启应用后,访问/actuator/health,就能看到名为kafkaProducerHealth的健康检查项。

方式二:将ProducerTemplate封装为Camel路由端点

把Kafka发送逻辑封装到Camel路由中,让Camel自动为其添加健康检查:

  1. 定义一个定时触发的健康检查路由:
import org.apache.camel.builder.RouteBuilder;
import org.springframework.stereotype.Component;

@Component
public class KafkaHealthCheckRoute extends RouteBuilder {

    @Override
    public void configure() throws Exception {
        // 每隔30秒发送一条测试消息到指定Kafka topic(需提前创建或配置自动创建)
        from("timer:kafka-health-check?period=30000")
                .routeId("kafka-health-check-route")
                .setBody(constant("kafka-health-check"))
                .to("kafka:health-check-topic?brokers=${kafka.bootstrap-servers}");
    }
}
  1. 此时Camel的健康检查会自动包含这个路由的Kafka生产者端点状态,你可以在/actuator/health的camelHealth节点下找到对应路由的健康信息。

注意事项

  • 确保你的依赖中包含camel-spring-boot-starter-health(Camel 3.x及以上版本,该依赖通常已包含在camel-spring-boot-starter中,无需额外引入)。
  • 如果使用测试消息的方式,要确保Kafka集群中存在对应的测试topic,或者配置Kafka自动创建topic。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 01:07:50