From a3c43647409ed731c5716d175d75630694f811f7 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} | 37 ++++++++++++++++--- .../kafka/kafka-binder/config-options.adoc | 6 +++ 4 files changed, 64 insertions(+), 6 deletions(-) rename binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/{KafkaBinderConfigurationPropertiesTest.java => KafkaBinderPropertiesTest.java} (76%) 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 95a668545..bc0386279 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 @@ -150,7 +150,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; /** * Earlier, @Autowired on this constructor was necessary for all the properties to be discovered @@ -504,6 +514,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. */ @@ -697,6 +715,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 76% 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 95b6bd2bc..1658c7992 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-2023 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,6 +23,8 @@ 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; @@ -32,9 +34,12 @@ 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; @@ -47,17 +52,39 @@ import static org.assertj.core.api.Assertions.assertThat; * @author Soby Chacko */ @SpringBootTest(classes = { KafkaBinderConfiguration.class, KafkaAutoConfiguration.class, - KafkaBinderConfigurationPropertiesTest.class }) -@TestPropertySource(locations = "classpath:binder-config.properties") -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 kafkaBinderConfigurationProperties() throws Exception { - assertThat(this.kafkaMessageChannelBinder).isNotNull(); + 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); diff --git a/docs/modules/ROOT/pages/kafka/kafka-binder/config-options.adoc b/docs/modules/ROOT/pages/kafka/kafka-binder/config-options.adoc index 3acbae24c..65c7a9a62 100644 --- a/docs/modules/ROOT/pages/kafka/kafka-binder/config-options.adoc +++ b/docs/modules/ROOT/pages/kafka/kafka-binder/config-options.adoc @@ -133,6 +133,12 @@ Enable Micrometer observation registry on all the bindings in this binder. + Default: false +spring.cloud.stream.kafka.binder.healthIndicatorConsumerGroup:: +`KafkaHealthIndicator` metadata consumer `group.id`. +This consumer is used by the `HealthIndicator` to query the metadata about the topics in use. ++ +Default: none. + [[kafka-consumer-properties]] == Kafka Consumer Properties