GH-2949: KafkaBinderHealthIndicator consumer group.id
* Add a new property in `KafkaBinderConfigurationProperties` to allow the users to specify a `group.id` for the metadata consumer used by the health indicator. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2949
This commit is contained in:
@@ -151,7 +151,17 @@ public class KafkaBinderConfigurationProperties {
|
||||
/**
|
||||
* Schema registry ssl configuration properties.
|
||||
*/
|
||||
private final String[] schemaRegistryProperties = new String[]{"schema.registry.url", "schema.registry.ssl.keystore.location", "schema.registry.ssl.keystore.password", "schema.registry.ssl.truststore.location", "schema.registry.ssl.truststore.password", "schema.registry.ssl.key.password"};
|
||||
private final String[] schemaRegistryProperties = new String[]{"schema.registry.url",
|
||||
"schema.registry.ssl.keystore.location", "schema.registry.ssl.keystore.password",
|
||||
"schema.registry.ssl.truststore.location", "schema.registry.ssl.truststore.password",
|
||||
"schema.registry.ssl.key.password"};
|
||||
|
||||
/**
|
||||
* Consumer group.id of the Kafka consumer in
|
||||
* {@link org.springframework.cloud.stream.binder.kafka.common.AbstractKafkaBinderHealthIndicator} that is used
|
||||
* for querying metadata from the broker (such as metadata information about the topics).
|
||||
*/
|
||||
private String healthIndicatorConsumerGroup;
|
||||
|
||||
/**
|
||||
* @Autowired on this constructor is necessary for all the properties to be discovered and bound when running as a native
|
||||
@@ -509,6 +519,14 @@ public class KafkaBinderConfigurationProperties {
|
||||
this.enableObservation = enableObservation;
|
||||
}
|
||||
|
||||
public String getHealthIndicatorConsumerGroup() {
|
||||
return healthIndicatorConsumerGroup;
|
||||
}
|
||||
|
||||
public void setHealthIndicatorConsumerGroup(String healthIndicatorConsumerGroup) {
|
||||
this.healthIndicatorConsumerGroup = healthIndicatorConsumerGroup;
|
||||
}
|
||||
|
||||
/**
|
||||
* Domain class that models transaction capabilities in Kafka.
|
||||
*/
|
||||
@@ -702,6 +720,8 @@ public class KafkaBinderConfigurationProperties {
|
||||
return this.kafkaProducerProperties;
|
||||
}
|
||||
|
||||
|
||||
|
||||
}
|
||||
|
||||
public static class Metrics {
|
||||
|
||||
@@ -34,6 +34,7 @@ import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.kafka.core.ConsumerFactory;
|
||||
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
|
||||
import org.springframework.util.ObjectUtils;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
* Configuration class for Kafka binder health indicator beans.
|
||||
@@ -67,6 +68,10 @@ public class KafkaBinderHealthIndicatorConfiguration {
|
||||
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
|
||||
configurationProperties.getKafkaConnectionString());
|
||||
}
|
||||
String healthIndicatorConsumerGroup = configurationProperties.getHealthIndicatorConsumerGroup();
|
||||
if (StringUtils.hasText(healthIndicatorConsumerGroup)) {
|
||||
props.put(ConsumerConfig.GROUP_ID_CONFIG, healthIndicatorConsumerGroup);
|
||||
}
|
||||
ConsumerFactory<?, ?> consumerFactory = new DefaultKafkaConsumerFactory<>(props);
|
||||
KafkaBinderHealthIndicator indicator = new KafkaBinderHealthIndicator(
|
||||
kafkaMessageChannelBinder, consumerFactory);
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2016-2022 the original author or authors.
|
||||
* Copyright 2016-2024 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -23,100 +23,125 @@ import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.apache.kafka.clients.consumer.ConsumerGroupMetadata;
|
||||
import org.apache.kafka.common.serialization.ByteArrayDeserializer;
|
||||
import org.apache.kafka.common.serialization.ByteArraySerializer;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.extension.ExtendWith;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
|
||||
import org.springframework.cloud.stream.binder.ExtendedProducerProperties;
|
||||
import org.springframework.cloud.stream.binder.kafka.common.AbstractKafkaBinderHealthIndicator;
|
||||
import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfiguration;
|
||||
import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties;
|
||||
import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties;
|
||||
import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties;
|
||||
import org.springframework.kafka.core.ConsumerFactory;
|
||||
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
|
||||
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
|
||||
import org.springframework.test.context.TestPropertySource;
|
||||
import org.springframework.test.context.junit.jupiter.SpringExtension;
|
||||
import org.springframework.util.ReflectionUtils;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
/**
|
||||
* @author Ilayaperumal Gopinathan
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
@ExtendWith(SpringExtension.class)
|
||||
@SpringBootTest(classes = { KafkaBinderConfiguration.class, KafkaAutoConfiguration.class,
|
||||
KafkaBinderConfigurationPropertiesTest.class })
|
||||
@TestPropertySource(locations = "classpath:binder-config.properties")
|
||||
public class KafkaBinderConfigurationPropertiesTest {
|
||||
KafkaBinderPropertiesTest.class })
|
||||
@TestPropertySource(locations = "classpath:binder-config.properties", properties =
|
||||
"spring.cloud.stream.kafka.binder.healthIndicatorConsumerGroup=health-consumer-group")
|
||||
class KafkaBinderPropertiesTest {
|
||||
|
||||
@Autowired
|
||||
private KafkaMessageChannelBinder kafkaMessageChannelBinder;
|
||||
|
||||
@Autowired
|
||||
private KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties;
|
||||
|
||||
@Autowired
|
||||
private KafkaBinderHealthIndicator kafkaBinderHealthIndicator;
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
void testKafkaBinderConfigurationProperties() throws Exception {
|
||||
assertThat(this.kafkaMessageChannelBinder).isNotNull();
|
||||
void kafkaBinderConfigurationProperties() throws Exception {
|
||||
assertThat(this.kafkaBinderConfigurationProperties).isNotNull();
|
||||
|
||||
// Testing a scenario in health indicator that is originally triggered by a property in KafkaBinderConfigurationProperties,
|
||||
// which ultimately creates a Kafka Consumer in the health indicator implementation.
|
||||
assertThat(this.kafkaBinderConfigurationProperties.getHealthIndicatorConsumerGroup())
|
||||
.isEqualTo("health-consumer-group");
|
||||
assertThat(this.kafkaBinderHealthIndicator).isNotNull();
|
||||
Field consumerFactoryField = AbstractKafkaBinderHealthIndicator.class.getDeclaredField("consumerFactory");
|
||||
consumerFactoryField.setAccessible(true);
|
||||
ConsumerFactory<?, ?> healthIndicatorConsumerFactory =
|
||||
(ConsumerFactory<?, ?>) consumerFactoryField.get(this.kafkaBinderHealthIndicator);
|
||||
assertThat(healthIndicatorConsumerFactory).isNotNull();
|
||||
Consumer<?, ?> consumer = healthIndicatorConsumerFactory.createConsumer();
|
||||
ConsumerGroupMetadata consumerGroupMetadata = consumer.groupMetadata();
|
||||
assertThat(consumerGroupMetadata.groupId()).isEqualTo("health-consumer-group");
|
||||
|
||||
KafkaProducerProperties kafkaProducerProperties = new KafkaProducerProperties();
|
||||
kafkaProducerProperties.setBufferSize(12345);
|
||||
kafkaProducerProperties.setBatchTimeout(100);
|
||||
kafkaProducerProperties.setCloseTimeout(10);
|
||||
kafkaProducerProperties
|
||||
.setCompressionType(KafkaProducerProperties.CompressionType.gzip);
|
||||
.setCompressionType(KafkaProducerProperties.CompressionType.gzip);
|
||||
ExtendedProducerProperties<KafkaProducerProperties> producerProperties = new ExtendedProducerProperties<>(
|
||||
kafkaProducerProperties);
|
||||
kafkaProducerProperties);
|
||||
Method getProducerFactoryMethod = KafkaMessageChannelBinder.class
|
||||
.getDeclaredMethod("getProducerFactory", String.class,
|
||||
ExtendedProducerProperties.class, String.class, String.class);
|
||||
.getDeclaredMethod("getProducerFactory", String.class,
|
||||
ExtendedProducerProperties.class, String.class, String.class);
|
||||
getProducerFactoryMethod.setAccessible(true);
|
||||
DefaultKafkaProducerFactory producerFactory = (DefaultKafkaProducerFactory) getProducerFactoryMethod
|
||||
.invoke(this.kafkaMessageChannelBinder, "bar", producerProperties, "bar.producer", "bar");
|
||||
.invoke(this.kafkaMessageChannelBinder, "bar", producerProperties, "bar.producer", "bar");
|
||||
Field producerFactoryConfigField = ReflectionUtils
|
||||
.findField(DefaultKafkaProducerFactory.class, "configs", Map.class);
|
||||
.findField(DefaultKafkaProducerFactory.class, "configs", Map.class);
|
||||
ReflectionUtils.makeAccessible(producerFactoryConfigField);
|
||||
Map<String, Object> producerConfigs = (Map<String, Object>) ReflectionUtils
|
||||
.getField(producerFactoryConfigField, producerFactory);
|
||||
.getField(producerFactoryConfigField, producerFactory);
|
||||
assertThat(producerConfigs.get("batch.size")).isEqualTo("12345");
|
||||
assertThat(producerConfigs.get("linger.ms")).isEqualTo("100");
|
||||
assertThat(producerConfigs.get("key.serializer"))
|
||||
.isEqualTo(ByteArraySerializer.class);
|
||||
.isEqualTo(ByteArraySerializer.class);
|
||||
assertThat(producerConfigs.get("value.serializer"))
|
||||
.isEqualTo(ByteArraySerializer.class);
|
||||
.isEqualTo(ByteArraySerializer.class);
|
||||
assertThat(producerConfigs.get("compression.type")).isEqualTo("gzip");
|
||||
|
||||
Field physicalCloseTimeoutField = ReflectionUtils
|
||||
.findField(DefaultKafkaProducerFactory.class, "physicalCloseTimeout", Duration.class);
|
||||
.findField(DefaultKafkaProducerFactory.class, "physicalCloseTimeout", Duration.class);
|
||||
ReflectionUtils.makeAccessible(physicalCloseTimeoutField);
|
||||
Duration physicalCloseTimeoutConfig = (Duration) ReflectionUtils
|
||||
.getField(physicalCloseTimeoutField, producerFactory);
|
||||
.getField(physicalCloseTimeoutField, producerFactory);
|
||||
assertThat(physicalCloseTimeoutConfig).isEqualTo(Duration.ofSeconds(10));
|
||||
|
||||
List<String> bootstrapServers = new ArrayList<>();
|
||||
bootstrapServers.add("10.98.09.199:9082");
|
||||
assertThat((((String) producerConfigs.get("bootstrap.servers"))
|
||||
.contains("10.98.09.199:9082"))).isTrue();
|
||||
.contains("10.98.09.199:9082"))).isTrue();
|
||||
Method createKafkaConsumerFactoryMethod = KafkaMessageChannelBinder.class
|
||||
.getDeclaredMethod("createKafkaConsumerFactory", boolean.class,
|
||||
String.class, ExtendedConsumerProperties.class, String.class, String.class);
|
||||
.getDeclaredMethod("createKafkaConsumerFactory", boolean.class,
|
||||
String.class, ExtendedConsumerProperties.class, String.class, String.class);
|
||||
createKafkaConsumerFactoryMethod.setAccessible(true);
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = new ExtendedConsumerProperties<>(
|
||||
new KafkaConsumerProperties());
|
||||
new KafkaConsumerProperties());
|
||||
DefaultKafkaConsumerFactory consumerFactory = (DefaultKafkaConsumerFactory) createKafkaConsumerFactoryMethod
|
||||
.invoke(this.kafkaMessageChannelBinder, true, "test", consumerProperties, "test.consumer", "test");
|
||||
.invoke(this.kafkaMessageChannelBinder, true, "test", consumerProperties, "test.consumer", "test");
|
||||
Field consumerFactoryConfigField = ReflectionUtils
|
||||
.findField(DefaultKafkaConsumerFactory.class, "configs", Map.class);
|
||||
.findField(DefaultKafkaConsumerFactory.class, "configs", Map.class);
|
||||
ReflectionUtils.makeAccessible(consumerFactoryConfigField);
|
||||
Map<String, Object> consumerConfigs = (Map<String, Object>) ReflectionUtils
|
||||
.getField(consumerFactoryConfigField, consumerFactory);
|
||||
.getField(consumerFactoryConfigField, consumerFactory);
|
||||
assertThat(consumerConfigs.get("key.deserializer"))
|
||||
.isEqualTo(ByteArrayDeserializer.class);
|
||||
.isEqualTo(ByteArrayDeserializer.class);
|
||||
assertThat(consumerConfigs.get("value.deserializer"))
|
||||
.isEqualTo(ByteArrayDeserializer.class);
|
||||
.isEqualTo(ByteArrayDeserializer.class);
|
||||
assertThat((((String) consumerConfigs.get("bootstrap.servers"))
|
||||
.contains("10.98.09.199:9082"))).isTrue();
|
||||
.contains("10.98.09.199:9082"))).isTrue();
|
||||
}
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user