From 7554ff9bbdc1690a22f14d93a8220e1e40913fe9 Mon Sep 17 00:00:00 2001 From: "Pommerening, Nico" Date: Wed, 22 Jun 2022 09:05:32 +0200 Subject: [PATCH] Kafka Streams binder health indicator improvements Fix KafkaStreamsBinderHealthIndicator overriding HealthCheck Thread Details to report full details. checkstyle fixes. --- .../KafkaStreamsBinderHealthIndicator.java | 20 +++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java index 6db6a64ce..16a5a534d 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java @@ -163,16 +163,20 @@ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator i if (isRunningResult) { final Set threadMetadata = kafkaStreams.metadataForLocalThreads(); + final Map threadDetails = new HashMap(); for (ThreadMetadata metadata : threadMetadata) { - 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 Map threadDetail = new HashMap(); + threadDetail.put("threadName", metadata.threadName()); + threadDetail.put("threadState", metadata.threadState()); + threadDetail.put("adminClientId", metadata.adminClientId()); + threadDetail.put("consumerClientId", metadata.consumerClientId()); + threadDetail.put("restoreConsumerClientId", metadata.restoreConsumerClientId()); + threadDetail.put("producerClientIds", metadata.producerClientIds()); + threadDetail.put("activeTasks", taskDetails(metadata.activeTasks())); + threadDetail.put("standbyTasks", taskDetails(metadata.standbyTasks())); + threadDetails.put(metadata.threadName(), threadDetail); } + perAppdIdDetails.put("threadDetails", threadDetails); final StreamsBuilderFactoryBean streamsBuilderFactoryBean = this.kafkaStreamsRegistry.streamBuilderFactoryBean(kafkaStreams); final String applicationId = (String) streamsBuilderFactoryBean.getStreamsConfiguration().get(StreamsConfig.APPLICATION_ID_CONFIG); details.put(applicationId, perAppdIdDetails);