From b833a9f371c57ea8a3bc3859f1bcfa31349fcb93 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 17 Oct 2019 19:24:53 -0400 Subject: [PATCH] Fix Kafka Streams binder health indicator issues When there are multiple Kafka Streams processors present, the health indicator overwrites the previous processor's health info. Addressing this issue. Resolves #771 --- .../KafkaStreamsBinderHealthIndicator.java | 30 ++++++++++++++----- .../kafka/streams/KafkaStreamsRegistry.java | 25 ++++++++++++++++ .../streams/StreamsBuilderFactoryManager.java | 6 +++- ...afkaStreamsBinderHealthIndicatorTests.java | 8 ++--- 4 files changed, 57 insertions(+), 12 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java index ca3292a97..19e195984 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java @@ -22,17 +22,20 @@ import java.util.Set; import java.util.stream.Collectors; import org.apache.kafka.streams.KafkaStreams; +import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.processor.TaskMetadata; import org.apache.kafka.streams.processor.ThreadMetadata; import org.springframework.boot.actuate.health.AbstractHealthIndicator; import org.springframework.boot.actuate.health.Health; import org.springframework.boot.actuate.health.Status; +import org.springframework.kafka.config.StreamsBuilderFactoryBean; /** * Health indicator for Kafka Streams. * * @author Arnaud Jardiné + * @author Soby Chacko */ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { @@ -53,15 +56,28 @@ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { builder.status(up ? Status.UP : Status.DOWN); } - private static Map buildDetails(KafkaStreams kStreams) { + private Map buildDetails(KafkaStreams kafkaStreams) { final Map details = new HashMap<>(); - if (kStreams.state().isRunning()) { - for (ThreadMetadata metadata : kStreams.localThreadsMetadata()) { - details.put("threadName", metadata.threadName()); - details.put("threadState", metadata.threadState()); - details.put("activeTasks", taskDetails(metadata.activeTasks())); - details.put("standbyTasks", taskDetails(metadata.standbyTasks())); + final Map perAppdIdDetails = new HashMap<>(); + if (kafkaStreams.state().isRunning()) { + for (ThreadMetadata metadata : kafkaStreams.localThreadsMetadata()) { + perAppdIdDetails.put("threadName", metadata.threadName()); + perAppdIdDetails.put("threadState", metadata.threadState()); + perAppdIdDetails.put("adminClientId", metadata.adminClientId()); + perAppdIdDetails.put("consumerClientId", metadata.consumerClientId()); + perAppdIdDetails.put("restoreConsumerClientId", metadata.restoreConsumerClientId()); + perAppdIdDetails.put("producerClientIds", metadata.producerClientIds()); + perAppdIdDetails.put("activeTasks", taskDetails(metadata.activeTasks())); + perAppdIdDetails.put("standbyTasks", taskDetails(metadata.standbyTasks())); } + final StreamsBuilderFactoryBean streamsBuilderFactoryBean = this.kafkaStreamsRegistry.streamBuilderFactoryBean(kafkaStreams); + final String applicationId = (String) streamsBuilderFactoryBean.getStreamsConfiguration().get(StreamsConfig.APPLICATION_ID_CONFIG); + details.put(applicationId, perAppdIdDetails); + } + else { + final StreamsBuilderFactoryBean streamsBuilderFactoryBean = this.kafkaStreamsRegistry.streamBuilderFactoryBean(kafkaStreams); + final String applicationId = (String) streamsBuilderFactoryBean.getStreamsConfiguration().get(StreamsConfig.APPLICATION_ID_CONFIG); + details.put(applicationId, String.format("The processor with application.id %s is down", applicationId)); } return details; } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java index 0799b89a9..88f4c28ab 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java @@ -16,11 +16,15 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.HashMap; import java.util.HashSet; +import java.util.Map; import java.util.Set; import org.apache.kafka.streams.KafkaStreams; +import org.springframework.kafka.config.StreamsBuilderFactoryBean; + /** * An internal registry for holding {@KafkaStreams} objects maintained through * {@link StreamsBuilderFactoryManager}. @@ -31,6 +35,8 @@ class KafkaStreamsRegistry { private final KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics; + private Map streamsStreamsBuilderFactoryBeanMap = new HashMap<>(); + KafkaStreamsRegistry(KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics) { this.kafkaStreamsBinderMetrics = kafkaStreamsBinderMetrics; } @@ -50,4 +56,23 @@ class KafkaStreamsRegistry { this.kafkaStreams.add(kafkaStreams); } + /** + * Make an association between {@link KafkaStreams} and its corresponding {@link StreamsBuilderFactoryBean}. + * + * @param kafkaStreams {@link KafkaStreams} object + * @param streamsBuilderFactoryBean Associtated {@link StreamsBuilderFactoryBean} for the {@link KafkaStreams} + */ + void addToStreamBuilderFactoryBeanMap(KafkaStreams kafkaStreams, StreamsBuilderFactoryBean streamsBuilderFactoryBean) { + streamsStreamsBuilderFactoryBeanMap.put(kafkaStreams, streamsBuilderFactoryBean); + } + + /** + * + * @param kafkaStreams {@link KafkaStreams} object + * @return Corresponding {@link StreamsBuilderFactoryBean}. + */ + StreamsBuilderFactoryBean streamBuilderFactoryBean(KafkaStreams kafkaStreams) { + return streamsStreamsBuilderFactoryBeanMap.get(kafkaStreams); + } + } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java index 7b775b6b0..92269b546 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java @@ -18,6 +18,8 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.util.Set; +import org.apache.kafka.streams.KafkaStreams; + import org.springframework.context.SmartLifecycle; import org.springframework.kafka.KafkaException; import org.springframework.kafka.config.StreamsBuilderFactoryBean; @@ -71,8 +73,10 @@ class StreamsBuilderFactoryManager implements SmartLifecycle { .getStreamsBuilderFactoryBeans(); for (StreamsBuilderFactoryBean streamsBuilderFactoryBean : streamsBuilderFactoryBeans) { streamsBuilderFactoryBean.start(); + final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); this.kafkaStreamsRegistry.registerKafkaStreams( - streamsBuilderFactoryBean.getKafkaStreams()); + kafkaStreams); + this.kafkaStreamsRegistry.addToStreamBuilderFactoryBeanMap(kafkaStreams, streamsBuilderFactoryBean); } this.running = true; } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java index ce54f25fd..541b76f06 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java @@ -77,7 +77,7 @@ public class KafkaStreamsBinderHealthIndicatorTests { @Test public void healthIndicatorUpTest() throws Exception { - try (ConfigurableApplicationContext context = singleStream()) { + try (ConfigurableApplicationContext context = singleStream("ApplicationHealthTest-xyz")) { receive(context, Lists.newArrayList(new ProducerRecord<>("in", "{\"id\":\"123\"}"), new ProducerRecord<>("in", "{\"id\":\"123\"}")), @@ -87,7 +87,7 @@ public class KafkaStreamsBinderHealthIndicatorTests { @Test public void healthIndicatorDownTest() throws Exception { - try (ConfigurableApplicationContext context = singleStream()) { + try (ConfigurableApplicationContext context = singleStream("ApplicationHealthTest-xyzabc")) { receive(context, Lists.newArrayList(new ProducerRecord<>("in", "{\"id\":\"123\"}"), new ProducerRecord<>("in", "{\"id\":\"124\"}")), @@ -186,7 +186,7 @@ public class KafkaStreamsBinderHealthIndicatorTests { assertThat(health.getStatus()).isEqualTo(expected); } - private ConfigurableApplicationContext singleStream() { + private ConfigurableApplicationContext singleStream(String applicationId) { SpringApplication app = new SpringApplication(KStreamApplication.class); app.setWebApplicationType(WebApplicationType.NONE); return app.run("--server.port=0", "--spring.jmx.enabled=false", @@ -198,7 +198,7 @@ public class KafkaStreamsBinderHealthIndicatorTests { "--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-xyz", + + applicationId, "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString()); }