Polish "Add Kafka health indicator"

Closes gh-11515
This commit is contained in:
Stephane Nicoll
2018-02-08 11:54:27 +01:00
parent 0dbd9429cc
commit 7cd19822c6
12 changed files with 114 additions and 108 deletions

View File

@@ -30,6 +30,7 @@ import org.apache.kafka.common.config.ConfigResource.Type;
import org.springframework.boot.actuate.health.AbstractHealthIndicator;
import org.springframework.boot.actuate.health.Health.Builder;
import org.springframework.boot.actuate.health.HealthIndicator;
import org.springframework.boot.actuate.health.Status;
import org.springframework.kafka.core.KafkaAdmin;
import org.springframework.util.Assert;
@@ -43,37 +44,35 @@ public class KafkaHealthIndicator extends AbstractHealthIndicator {
static final String REPLICATION_PROPERTY = "transaction.state.log.replication.factor";
private final KafkaAdmin kafkaAdmin;
private final DescribeClusterOptions describeOptions;
/**
* Create a new {@link KafkaHealthIndicator} instance.
*
* @param kafkaAdmin the kafka admin
* @param responseTimeout the describe cluster request timeout in milliseconds
* @param requestTimeout the request timeout in milliseconds
*/
public KafkaHealthIndicator(KafkaAdmin kafkaAdmin, long responseTimeout) {
public KafkaHealthIndicator(KafkaAdmin kafkaAdmin, long requestTimeout) {
Assert.notNull(kafkaAdmin, "KafkaAdmin must not be null");
this.kafkaAdmin = kafkaAdmin;
this.describeOptions = new DescribeClusterOptions()
.timeoutMs((int) responseTimeout);
.timeoutMs((int) requestTimeout);
}
@Override
protected void doHealthCheck(Builder builder) throws Exception {
try (AdminClient adminClient = AdminClient.create(this.kafkaAdmin.getConfig())) {
DescribeClusterResult result = adminClient.describeCluster(this.describeOptions);
DescribeClusterResult result = adminClient.describeCluster(
this.describeOptions);
String brokerId = result.controller().get().idString();
int replicationFactor = getReplicationFactor(brokerId, adminClient);
int nodes = result.nodes().get().size();
if (nodes >= replicationFactor) {
builder.up();
}
else {
builder.down();
}
builder.withDetail("clusterId", result.clusterId().get());
builder.withDetail("brokerId", brokerId);
builder.withDetail("nodes", nodes);
Status status = nodes >= replicationFactor ? Status.UP : Status.DOWN;
builder.status(status)
.withDetail("clusterId", result.clusterId().get())
.withDetail("brokerId", brokerId)
.withDetail("nodes", nodes);
}
}
@@ -85,5 +84,6 @@ public class KafkaHealthIndicator extends AbstractHealthIndicator {
Config brokerConfig = kafkaConfig.get(configResource);
return Integer.parseInt(brokerConfig.get(REPLICATION_PROPERTY).value());
}
}

View File

@@ -15,6 +15,6 @@
*/
/**
* Actuator support for Kafka.
* Actuator support for Apache Kafka.
*/
package org.springframework.boot.actuate.kafka;

View File

@@ -20,27 +20,73 @@ import java.util.Collections;
import java.util.Map;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.junit.After;
import org.junit.Test;
import org.springframework.boot.actuate.health.Health;
import org.springframework.boot.actuate.health.Status;
import org.springframework.kafka.core.KafkaAdmin;
import org.springframework.kafka.test.rule.KafkaEmbedded;
import org.springframework.util.SocketUtils;
import static org.assertj.core.api.Assertions.assertThat;
/**
* Test for {@link KafkaHealthIndicator}
* Tests for {@link KafkaHealthIndicator}.
*
* @author Juan Rada
* @author Stephane Nicoll
*/
public class KafkaHealthIndicatorTests {
private static final Long RESPONSE_TIME = 1000L;
private KafkaEmbedded kafkaEmbedded;
private KafkaAdmin kafkaAdmin;
@After
public void shutdownKafka() throws Exception {
if (this.kafkaEmbedded != null) {
this.kafkaEmbedded.destroy();
}
}
@Test
public void kafkaIsUp() throws Exception {
startKafka(1);
KafkaHealthIndicator healthIndicator =
new KafkaHealthIndicator(this.kafkaAdmin, 1000L);
Health health = healthIndicator.health();
assertThat(health.getStatus()).isEqualTo(Status.UP);
assertDetails(health.getDetails());
}
@Test
public void kafkaIsDown() {
int freePort = SocketUtils.findAvailableTcpPort();
this.kafkaAdmin = new KafkaAdmin(Collections.singletonMap(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:" + freePort));
KafkaHealthIndicator healthIndicator =
new KafkaHealthIndicator(this.kafkaAdmin, 1L);
Health health = healthIndicator.health();
assertThat(health.getStatus()).isEqualTo(Status.DOWN);
assertThat((String) health.getDetails().get("error")).isNotEmpty();
}
@Test
public void notEnoughNodesForReplicationFactor() throws Exception {
startKafka(2);
KafkaHealthIndicator healthIndicator =
new KafkaHealthIndicator(this.kafkaAdmin, 1000L);
Health health = healthIndicator.health();
assertThat(health.getStatus()).isEqualTo(Status.DOWN);
assertDetails(health.getDetails());
}
private void assertDetails(Map<String, Object> details) {
assertThat(details).containsEntry("brokerId", "0");
assertThat(details).containsKey("clusterId");
assertThat(details).containsEntry("nodes", 1);
}
private void startKafka(int replicationFactor) throws Exception {
this.kafkaEmbedded = new KafkaEmbedded(1, true);
this.kafkaEmbedded.brokerProperties(Collections.singletonMap(
@@ -52,46 +98,4 @@ public class KafkaHealthIndicatorTests {
this.kafkaEmbedded.getBrokersAsString()));
}
private void shutdownKafka() throws Exception {
this.kafkaEmbedded.destroy();
}
@Test
public void kafkaIsUp() throws Exception {
startKafka(1);
KafkaHealthIndicator healthIndicator =
new KafkaHealthIndicator(this.kafkaAdmin, RESPONSE_TIME);
Health health = healthIndicator.health();
assertThat(health.getStatus()).isEqualTo(Status.UP);
assertDetails(health.getDetails());
shutdownKafka();
}
private void assertDetails(Map<String, Object> details) {
assertThat(details).containsEntry("brokerId", "0");
assertThat(details).containsKey("clusterId");
assertThat(details).containsEntry("nodes", 1);
}
@Test
public void notEnoughNodesForReplicationFactor() throws Exception {
startKafka(2);
KafkaHealthIndicator healthIndicator =
new KafkaHealthIndicator(this.kafkaAdmin, RESPONSE_TIME);
Health health = healthIndicator.health();
assertThat(health.getStatus()).isEqualTo(Status.DOWN);
assertDetails(health.getDetails());
shutdownKafka();
}
@Test
public void kafkaIsDown() throws Exception {
this.kafkaAdmin = new KafkaAdmin(Collections.singletonMap(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:34987"));
KafkaHealthIndicator healthIndicator =
new KafkaHealthIndicator(this.kafkaAdmin, RESPONSE_TIME);
Health health = healthIndicator.health();
assertThat(health.getStatus()).isEqualTo(Status.DOWN);
assertThat((String) health.getDetails().get("error")).isNotEmpty();
}
}