From 720e025f4e175642e12938a2dc73660fe328ba3a Mon Sep 17 00:00:00 2001 From: 98001yash Date: Thu, 8 Oct 2026 23:01:40 +0530 Subject: [PATCH] Add Kafka Streams health check for unreachable broker Signed-off-by: 98001yash --- ...afkaStreamsBinderHealthIndicatorTests.java | 37 +++++++++++++++++++ 1 file changed, 37 insertions(+) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java index 40afaeb65b..516d0f6199 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java @@ -81,6 +81,22 @@ void healthIndicatorUpTest() throws Exception { } } + @Test + void healthIndicatorDownWhenKafkaBrokerIsNotReachable() throws Exception { + try (ConfigurableApplicationContext context = singleStreamWithUnreachableBroker()) { + TimeUnit.SECONDS.sleep(2); + + CompositeHealthContributor healthIndicator = context + .getBean("bindersHealthContributor", CompositeHealthContributor.class); + KafkaStreamsBinderHealthIndicator kafkaStreamsBinderHealthIndicator = + (KafkaStreamsBinderHealthIndicator) healthIndicator.getContributor("kstream"); + + Health health = kafkaStreamsBinderHealthIndicator.health(); + + assertThat(health.getStatus()).isEqualTo(Status.DOWN); + } + } + @Test void healthIndicatorUpMultipleCallsTest() throws Exception { try (ConfigurableApplicationContext context = singleStream("ApplicationHealthTest-xyz")) { @@ -208,6 +224,27 @@ private ConfigurableApplicationContext singleStream(String applicationId) { + embeddedKafka.getBrokersAsString()); } + private ConfigurableApplicationContext singleStreamWithUnreachableBroker() { + SpringApplication app = new SpringApplication(KStreamApplication.class); + app.setWebApplicationType(WebApplicationType.NONE); + + return app.run( + "--server.port=0", + "--spring.jmx.enabled=false", + "--spring.cloud.stream.function.bindings.process-in-0=input", + "--spring.cloud.stream.function.bindings.process-out-0=output", + "--spring.cloud.stream.bindings.input.destination=in", + "--spring.cloud.stream.bindings.output.destination=out", + "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde=" + + "org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=" + + "org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=" + + "ApplicationHealthTest-unreachable", + "--spring.cloud.stream.kafka.binder.healthTimeout=1", + "--spring.cloud.stream.kafka.streams.binder.brokers=localhost:65535"); + } + private ConfigurableApplicationContext multipleStream() { System.setProperty("logging.level.org.apache.kafka", "OFF"); SpringApplication app = new SpringApplication(AnotherKStreamApplication.class);