Tengo algo como a continuación que funciona bien, pero preferiría verificar el estado sin enviar ningún mensaje (no solo verificar la conexión del socket). Sé que Kafka tiene algo como KafkaHealthIndicator listo para usar, ¿alguien tiene experiencia o ejemplo al usarlo?
public class KafkaHealthIndicator implements HealthIndicator { private final Logger log = LoggerFactory.getLogger(KafkaHealthIndicator.class); private KafkaTemplate<String, String> kafka; public KafkaHealthIndicator(KafkaTemplate<String, String> kafka) { this.kafka = kafka; } @Override public Health health() { try { kafka.send("kafka-health-indicator", "❥").get(100, TimeUnit.MILLISECONDS); } catch (InterruptedException | ExecutionException | TimeoutException e) { return Health.down(e).build(); } return Health.up().build(); } }kafkaAdminClient.describeCluster(..) es el punto donde se prueba la disponibilidad de Kafka.
@Configuration public class KafkaConfig { @Autowired private KafkaAdmin kafkaAdmin; @Bean public AdminClient kafkaAdminClient() { return AdminClient.create(kafkaAdmin.getConfigurationProperties()); } @Bean public HealthIndicator kafkaHealthIndicator(AdminClient kafkaAdminClient) { final DescribeClusterOptions options = new DescribeClusterOptions() .timeoutMs(1000); return new AbstractHealthIndicator() { @Override protected void doHealthCheck(Health.Builder builder) throws Exception { // When Kafka is not connected, describeCluster() method throws // an exception which in turn sets this indicator as being DOWN. kafkaAdminClient.describeCluster(options); builder.up().build(); } }; } }Para una sonda más detallada, agregue:
DescribeClusterResult clusterDesc = kafkaAdminClient.describeCluster(options); builder.up() .withDetail("clusterId", clusterDesc.clusterId().get()) .withDetail("nodeCount", clusterDesc.nodes().get().size()) .build();Utilice la API de AdminClient para comprobar el estado del clúster describiendo el clúster y/o los temas con los que interactuará y verificando que esos temas tengan la cantidad necesaria de réplicas sincronizadas, por ejemplo.
Kafka tiene algo como KafkaHealthIndicator listo para usar
no lo hace La integración de Spring con Kafka podría