From d76f0bc42e69d9732fcd4d11e5ada2364a154442 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 13 May 2024 20:41:24 -0400 Subject: [PATCH] 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 --- .../KafkaBinderConfigurationProperties.java | 22 ++++- ...fkaBinderHealthIndicatorConfiguration.java | 5 ++ ...st.java => KafkaBinderPropertiesTest.java} | 85 ++++++++++++------- 3 files changed, 81 insertions(+), 31 deletions(-) rename binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/{KafkaBinderConfigurationPropertiesTest.java => KafkaBinderPropertiesTest.java} (58%) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java index c06cef0ea..6fe00135e 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java @@ -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 { diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderHealthIndicatorConfiguration.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderHealthIndicatorConfiguration.java index 67d53a182..53104a5cc 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderHealthIndicatorConfiguration.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderHealthIndicatorConfiguration.java @@ -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); diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderPropertiesTest.java similarity index 58% rename from binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java rename to binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderPropertiesTest.java index dbafc94fd..456808e91 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderPropertiesTest.java @@ -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 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 producerConfigs = (Map) 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 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 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 consumerConfigs = (Map) 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(); } }