From b1b0e960d78d77c830fb02acbdbe88d476b3a25f Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Fri, 19 Feb 2016 16:48:33 -0500 Subject: [PATCH] Corrections for health indicator Ensure that try-catch is applied around connecting to Zookeeper as well Register topics in use for consumer --- .../kafka/KafkaBinderHealthIndicator.java | 33 ++++++++++--------- .../kafka/KafkaMessageChannelBinder.java | 2 +- 2 files changed, 18 insertions(+), 17 deletions(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java index b07986fa6..75dea05c2 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java @@ -48,15 +48,29 @@ public class KafkaBinderHealthIndicator implements HealthIndicator { @Override public Health health() { - ZkClient zkClient = new ZkClient(binder.getZkAddress(), binder.getZkSessionTimeout(), - binder.getZkConnectionTimeout(), ZKStringSerializer$.MODULE$); - Set brokersInClusterSet = new HashSet<>(); + ZkClient zkClient = null; try { + zkClient = new ZkClient(binder.getZkAddress(), binder.getZkSessionTimeout(), + binder.getZkConnectionTimeout(), ZKStringSerializer$.MODULE$); + Set brokersInClusterSet = new HashSet<>(); Seq allBrokersInCluster = ZkUtils$.MODULE$.getAllBrokersInCluster(zkClient); Collection brokersInCluster = JavaConversions.asJavaCollection(allBrokersInCluster); for (Broker broker : brokersInCluster) { brokersInClusterSet.add(broker.connectionString()); } + Set downMessages = new HashSet<>(); + for (Map.Entry> entry : binder.getTopicsInUse().entrySet()) { + for (Partition partition : entry.getValue()) { + BrokerAddress address = binder.getConnectionFactory().getLeader(partition); + if (!brokersInClusterSet.contains(address.toString())) { + downMessages.add(address.toString()); + } + } + } + if (downMessages.isEmpty()) { + return Health.up().build(); + } + return Health.down().withDetail("Following brokers are down: ", downMessages.toString()).build(); } catch (Exception e) { return Health.down(e).build(); @@ -71,18 +85,5 @@ public class KafkaBinderHealthIndicator implements HealthIndicator { } } } - Set downMessages = new HashSet<>(); - for (Map.Entry> entry : binder.getTopicsInUse().entrySet()) { - for (Partition partition : entry.getValue()) { - BrokerAddress address = binder.getConnectionFactory().getLeader(partition); - if (!brokersInClusterSet.contains(address.toString())) { - downMessages.add(address.toString()); - } - } - } - if (downMessages.isEmpty()) { - return Health.up().build(); - } - return Health.down().withDetail("Following brokers are down: ", downMessages.toString()).build(); } } diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 132a1b140..031b943e6 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -600,7 +600,7 @@ public class KafkaMessageChannelBinder extends MessageChannelBinderSupport { } } } - + topicsInUse.put(name, listenedPartitions); ReceivingHandler rh = new ReceivingHandler(); rh.setOutputChannel(moduleInputChannel);