From f5008839947699a67e13401c0c11210b47b7bbd0 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 12 Dec 2022 15:56:01 -0500 Subject: [PATCH] Observation related changes in Kafka binder (#2582) * Observation related changes in Kafka binder Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2576 * Follow up to the previous commit on observation changes * Address PR review * Addressing PR review * Addressing PR review --- .../KafkaBinderConfigurationProperties.java | 13 +++++ .../spring-cloud-stream-binder-kafka/pom.xml | 5 ++ .../kafka/KafkaMessageChannelBinder.java | 18 ++++++ .../stream/binder/kafka/KafkaBinderTests.java | 56 +++++++++++++++++++ .../binder/AbstractMessageChannelBinder.java | 2 +- .../main/asciidoc/kafka/kafka_overview.adoc | 6 ++ 6 files changed, 99 insertions(+), 1 deletion(-) 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 969b805ac..91d36ae69 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 @@ -139,6 +139,11 @@ public class KafkaBinderConfigurationProperties { */ private String certificateStoreDirectory; + /** + * Enable Micrometer observation registry across all the bindings in the binder. + */ + private boolean enableObservation; + public KafkaBinderConfigurationProperties(KafkaProperties kafkaProperties) { Assert.notNull(kafkaProperties, "'kafkaProperties' cannot be null"); this.kafkaProperties = kafkaProperties; @@ -477,6 +482,14 @@ public class KafkaBinderConfigurationProperties { this.certificateStoreDirectory = certificateStoreDirectory; } + public boolean isEnableObservation() { + return this.enableObservation; + } + + public void setEnableObservation(boolean enableObservation) { + this.enableObservation = enableObservation; + } + /** * Domain class that models transaction capabilities in Kafka. */ diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/pom.xml b/binders/kafka-binder/spring-cloud-stream-binder-kafka/pom.xml index 49864bf9d..0f5883a5d 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/pom.xml +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/pom.xml @@ -61,6 +61,11 @@ awaitility test + + io.micrometer + micrometer-observation-test + test + diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 455669284..ed3e5a175 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -51,6 +51,7 @@ import org.apache.kafka.common.header.internals.RecordHeader; import org.apache.kafka.common.header.internals.RecordHeaders; import org.springframework.beans.factory.DisposableBean; +import org.springframework.beans.factory.SmartInitializingSingleton; import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder; import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.binder.BinderSpecificPropertiesProvider; @@ -502,9 +503,21 @@ public class KafkaMessageChannelBinder extends } handler.setHeaderMapper(mapper); + + if (this.configurationProperties.isEnableObservation()) { + kafkaTemplate.setObservationEnabled(true); + } + kafkaTemplate.setApplicationContext(getApplicationContext()); + return handler; } + @Override + @SuppressWarnings("rawtypes") + protected void customizeProducerMessageHandler(MessageHandler producerMessageHandler, String destinationName) { + super.customizeProducerMessageHandler(producerMessageHandler, destinationName); + ((KafkaProducerMessageHandler) producerMessageHandler).getKafkaTemplate().afterSingletonsInstantiated(); + } @Override protected void postProcessOutputChannel(MessageChannel outputChannel, @@ -619,6 +632,11 @@ public class KafkaMessageChannelBinder extends ? new ContainerProperties(Pattern.compile(topics[0])) : new ContainerProperties(topics) : new ContainerProperties(topicPartitionOffsets); + + if (this.configurationProperties.isEnableObservation()) { + containerProperties.setObservationEnabled(true); + } + KafkaAwareTransactionManager transMan = transactionManager( extendedConsumerProperties.getExtension().getTransactionManager()); if (transMan != null) { diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 150cc26c1..f819b62d3 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -38,6 +38,8 @@ import java.util.concurrent.atomic.AtomicReference; import java.util.stream.IntStream; import com.fasterxml.jackson.databind.ObjectMapper; +import io.micrometer.observation.ObservationRegistry; +import io.micrometer.observation.tck.TestObservationRegistry; import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.CreateTopicsResult; @@ -98,6 +100,7 @@ import org.springframework.cloud.stream.binder.kafka.utils.DlqPartitionFunction; import org.springframework.cloud.stream.binder.kafka.utils.KafkaTopicUtils; import org.springframework.cloud.stream.binding.MessageConverterConfigurer.PartitioningInterceptor; import org.springframework.cloud.stream.config.BindingProperties; +import org.springframework.cloud.stream.config.ProducerMessageHandlerCustomizer; import org.springframework.cloud.stream.provisioning.ProvisioningException; import org.springframework.context.ApplicationContext; import org.springframework.context.ConfigurableApplicationContext; @@ -3889,6 +3892,59 @@ public class KafkaBinderTests extends .withCauseExactlyInstanceOf(IllegalStateException.class); } + @Test + void testObservationEnabledOnTheBinder() throws Exception { + KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties = createConfigurationProperties(); + kafkaBinderConfigurationProperties.setEnableObservation(true); + AbstractKafkaTestBinder binder = getBinder(kafkaBinderConfigurationProperties); + + setupBindingAndAssert("enable-observation.1", binder); + } + + @SuppressWarnings("rawtypes") + @Test + void testObservationEnabledThroughProducerMessageHandlerCustomizer() throws Exception { + AbstractKafkaTestBinder binder = getBinder(); + KafkaMessageChannelBinder kafkaMessageChannelBinder = binder.getCoreBinder(); + kafkaMessageChannelBinder.setProducerMessageHandlerCustomizer( + (ProducerMessageHandlerCustomizer) (handler, destinationName) -> + handler.getKafkaTemplate().setObservationEnabled(true)); + + setupBindingAndAssert("enable-observation.2", binder); + } + + private void setupBindingAndAssert(String bindingName, AbstractKafkaTestBinder binder) throws Exception { + ConfigurableApplicationContext applicationContext = (ConfigurableApplicationContext) binder.getApplicationContext(); + TestObservationRegistry observationRegistry = TestObservationRegistry.create(); + + applicationContext.getBeanFactory().registerSingleton("test-registry", observationRegistry); + + DirectChannel moduleOutputChannel = createBindableChannel("output", + new BindingProperties()); + ExtendedProducerProperties producerProps = new ExtendedProducerProperties<>( + new KafkaProducerProperties()); + Binding producerBinding = binder.bindProducer(bindingName, + moduleOutputChannel, producerProps); + + assertionsOnKafkaTemplate(observationRegistry, producerBinding); + } + + @SuppressWarnings("rawtypes") + private static void assertionsOnKafkaTemplate(TestObservationRegistry observationRegistry, Binding producerBinding) { + KafkaProducerMessageHandler endpoint = TestUtils.getPropertyValue(producerBinding, + "lifecycle", KafkaProducerMessageHandler.class); + + final KafkaTemplate kafkaTemplate = (KafkaTemplate) new DirectFieldAccessor(endpoint).getPropertyValue("kafkaTemplate"); + assertThat(kafkaTemplate).isNotNull(); + Boolean observationEnabled = (Boolean) new DirectFieldAccessor(kafkaTemplate).getPropertyValue("observationEnabled"); + assertThat(observationEnabled).isTrue(); + + ObservationRegistry observationRegistry1 = (ObservationRegistry) new DirectFieldAccessor(kafkaTemplate).getPropertyValue("observationRegistry"); + assertThat(observationRegistry).isSameAs(observationRegistry1); + + producerBinding.unbind(); + } + private final class FailingInvocationCountingMessageHandler implements MessageHandler { diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index 252e087f1..014022c50 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -355,7 +355,7 @@ public abstract class AbstractMessageChannelBinder