From 59c170d1472a48f6d644174fcd29944e735c19e9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Flaviu=20Mure=C8=99an?= Date: Fri, 7 May 2021 22:20:37 +0200 Subject: [PATCH] Kafka instrumentation enhancements (#1936) * Bump javadocs since version to 3.1.0 * Extend integration tests to cover all Kafka clients * Add autoconfiguration module for kafka instrumentation * Refactor instrumentation for reactive Kafka Receiver * Refactor instrumentation for reactive Kafka Receiver * Add docs for Kafka instrumentation * Revert "Refactor instrumentation for reactive Kafka Receiver" This reverts commit 58c8f2fa * Revert "Revert "Refactor instrumentation for reactive Kafka Receiver"" This reverts commit 450a9f8c * Remove empty test * Resolve comments from PR 1936 * Revert whitespaces * Revert whitespaces in common tests pom.xml * Fix autoconfig to consider generics when registering beans. Only register reactor-kafka beans if the dependency is on the classpath. * Split autoconfig for kafka and reactor-kafka --- .../main/asciidoc/documentation-overview.adoc | 3 +- docs/src/main/asciidoc/integrations.adoc | 19 ++ spring-cloud-sleuth-autoconfigure/pom.xml | 5 + .../kafka/TracingKafkaAutoConfiguration.java | 74 ++++++++ ...TracingKafkaConsumerBeanPostProcessor.java | 49 +++++ ...TracingKafkaProducerBeanPostProcessor.java | 49 +++++ .../TracingReactorKafkaAutoConfiguration.java | 61 +++++++ ...TraceSpringMessagingAutoConfiguration.java | 4 +- ...itional-spring-configuration-metadata.json | 6 + .../main/resources/META-INF/spring.factories | 1 + .../kafka/KafkaTracingCallback.java | 2 +- .../kafka/TracingKafkaConsumer.java | 39 +++- .../kafka/TracingKafkaConsumerFactory.java | 49 +++++ .../kafka/TracingKafkaProducer.java | 53 ++++-- .../kafka/TracingKafkaProducerFactory.java | 17 +- .../kafka/TracingKafkaPropagatorGetter.java | 2 +- .../kafka/TracingKafkaPropagatorSetter.java | 2 +- .../kafka/TracingKafkaReceiver.java | 112 ------------ .../kafka/TracingKafkaConsumerTest.java | 14 +- .../kafka/TracingKafkaProducerTest.java | 17 +- .../kafka/TracingKafkaReceiverTest.java | 61 ------- .../pom.xml | 5 +- .../instrument/kafka/KafkaConsumerTest.java | 65 +++++++ .../instrument/kafka/KafkaProducerTest.java | 28 +++ .../instrument/kafka/KafkaReceiverTest.java | 64 +++++++ .../instrument/kafka/KafkaSenderTest.java | 69 +++++++ tests/common/pom.xml | 5 + .../instrument/kafka/KafkaConsumerTest.java | 159 ++++++++++++++++ .../instrument/kafka/KafkaProducerTest.java | 93 ++++++++-- .../instrument/kafka/KafkaReceiverTest.java | 150 ++++++++++++++++ .../instrument/kafka/KafkaSenderTest.java | 169 ++++++++++++++++++ .../instrument/kafka/KafkaTestUtils.java | 52 ++++++ 32 files changed, 1266 insertions(+), 232 deletions(-) create mode 100644 spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaAutoConfiguration.java create mode 100644 spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaConsumerBeanPostProcessor.java create mode 100644 spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaProducerBeanPostProcessor.java create mode 100644 spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingReactorKafkaAutoConfiguration.java create mode 100644 spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumerFactory.java delete mode 100644 spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaReceiver.java delete mode 100644 spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaReceiverTest.java create mode 100644 tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaConsumerTest.java create mode 100644 tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaReceiverTest.java create mode 100644 tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaSenderTest.java create mode 100644 tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaConsumerTest.java create mode 100644 tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaReceiverTest.java create mode 100644 tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaSenderTest.java create mode 100644 tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTestUtils.java diff --git a/docs/src/main/asciidoc/documentation-overview.adoc b/docs/src/main/asciidoc/documentation-overview.adoc index 9bae80a47..741e3bcca 100644 --- a/docs/src/main/asciidoc/documentation-overview.adoc +++ b/docs/src/main/asciidoc/documentation-overview.adoc @@ -77,6 +77,7 @@ Need more details about {project-full-name}'s core features? Finally, we have topics related to instrumentation integrations: * *Integrations:* +<> | <> | <> | <> | @@ -91,4 +92,4 @@ Finally, we have topics related to instrumentation integrations: <> | <> | <> | -<> \ No newline at end of file +<> diff --git a/docs/src/main/asciidoc/integrations.adoc b/docs/src/main/asciidoc/integrations.adoc index c710467bc..12670a2b0 100644 --- a/docs/src/main/asciidoc/integrations.adoc +++ b/docs/src/main/asciidoc/integrations.adoc @@ -5,6 +5,25 @@ include::_attributes.adoc[] In this section, we describe how to customize various parts of Spring Cloud Sleuth. +[[sleuth-kafka-integration]] +== Apache Kafka + +This feature is available for all tracer implementations. + +We decorate the Kafka clients (`KafkaProducer` and `KafkaConsumer`) to create a span for each event that is produced or consumed. You can disable this feature by setting the value of `spring.sleuth.kafka.enabled` to `false`. + +IMPORTANT: You have to register the `Producer` or `Consumer` as beans in order for Sleuth's auto-configuration to decorate them. When you then inject the beans, the expected type must be `Producer` or `Consumer` (and NOT e.g. `KafkaProducer`). + +We also provide `TracingKafkaProducerFactory` and `TracingKafkaConsumerFactory` to be used with the https://projectreactor.io/docs/kafka/release/reference/[Reactor Kafka] clients (`KafkaSender` and `KafkaReceiver`, respectively). See an example in the snippet below: + +[source,java,indent=0] +---- +@Bean +KafkaReceiver reactiveKafkaReceiver(TracingKafkaConsumerFactory tracingKafkaConsumerFactory, KafkaReceiverOptions kafkaReceiverOptions) { + return KafkaReceiver.create(tracingKafkaConsumerFactory, kafkaReceiverOptions); +} +---- + [[sleuth-async-integration]] == Asynchronous Communication diff --git a/spring-cloud-sleuth-autoconfigure/pom.xml b/spring-cloud-sleuth-autoconfigure/pom.xml index 1c74bedd6..970d8d5a7 100644 --- a/spring-cloud-sleuth-autoconfigure/pom.xml +++ b/spring-cloud-sleuth-autoconfigure/pom.xml @@ -275,6 +275,11 @@ spring-jms true + + io.projectreactor.kafka + reactor-kafka + true + io.github.lognet diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaAutoConfiguration.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaAutoConfiguration.java new file mode 100644 index 000000000..4eba0bf43 --- /dev/null +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaAutoConfiguration.java @@ -0,0 +1,74 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.autoconfig.instrument.kafka; + +import org.apache.kafka.clients.KafkaClient; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.producer.ProducerRecord; + +import org.springframework.beans.factory.BeanFactory; +import org.springframework.boot.autoconfigure.AutoConfigureAfter; +import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.autoconfig.brave.BraveAutoConfiguration; +import org.springframework.cloud.sleuth.instrument.kafka.TracingKafkaPropagatorGetter; +import org.springframework.cloud.sleuth.instrument.kafka.TracingKafkaPropagatorSetter; +import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +/** + * {@link org.springframework.boot.autoconfigure.EnableAutoConfiguration + * Auto-configuration} that registers instrumentation for Kafka. + * + * @author Anders Clausen + * @author Flaviu Muresan + * @since 3.1.0 + */ +@Configuration(proxyBeanMethods = false) +@ConditionalOnClass(KafkaClient.class) +@ConditionalOnBean(Tracer.class) +@AutoConfigureAfter(BraveAutoConfiguration.class) +@ConditionalOnProperty(value = "spring.sleuth.kafka.enabled", matchIfMissing = true) +public class TracingKafkaAutoConfiguration { + + @Bean + @ConditionalOnMissingBean(value = ProducerRecord.class, parameterizedContainer = Propagator.Setter.class) + Propagator.Setter> tracingKafkaPropagationSetter() { + return new TracingKafkaPropagatorSetter(); + } + + @Bean + @ConditionalOnMissingBean(value = ConsumerRecord.class, parameterizedContainer = Propagator.Getter.class) + Propagator.Getter> tracingKafkaPropagationGetter() { + return new TracingKafkaPropagatorGetter(); + } + + @Bean + static TracingKafkaProducerBeanPostProcessor tracingKafkaProducerBeanPostProcessor(BeanFactory beanFactory) { + return new TracingKafkaProducerBeanPostProcessor(beanFactory); + } + + @Bean + static TracingKafkaConsumerBeanPostProcessor tracingKafkaConsumerBeanPostProcessor(BeanFactory beanFactory) { + return new TracingKafkaConsumerBeanPostProcessor(beanFactory); + } + +} diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaConsumerBeanPostProcessor.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaConsumerBeanPostProcessor.java new file mode 100644 index 000000000..dbcd976a4 --- /dev/null +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaConsumerBeanPostProcessor.java @@ -0,0 +1,49 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.autoconfig.instrument.kafka; + +import org.apache.kafka.clients.consumer.Consumer; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.config.BeanPostProcessor; +import org.springframework.cloud.sleuth.instrument.kafka.TracingKafkaConsumer; + +/** + * Bean post processor for {@link org.apache.kafka.clients.consumer.Consumer}. + * + * @author Anders Clausen + * @author Flaviu Muresan + * @since 3.1.0 + */ +public class TracingKafkaConsumerBeanPostProcessor implements BeanPostProcessor { + + private final BeanFactory beanFactory; + + public TracingKafkaConsumerBeanPostProcessor(BeanFactory beanFactory) { + this.beanFactory = beanFactory; + } + + @Override + public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { + if (bean instanceof Consumer && !(bean instanceof TracingKafkaConsumer)) { + return new TracingKafkaConsumer<>((Consumer) bean, beanFactory); + } + return bean; + } + +} diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaProducerBeanPostProcessor.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaProducerBeanPostProcessor.java new file mode 100644 index 000000000..9ff34bd3e --- /dev/null +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingKafkaProducerBeanPostProcessor.java @@ -0,0 +1,49 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.autoconfig.instrument.kafka; + +import org.apache.kafka.clients.producer.Producer; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.config.BeanPostProcessor; +import org.springframework.cloud.sleuth.instrument.kafka.TracingKafkaProducer; + +/** + * Bean post processor for {@link org.apache.kafka.clients.producer.Producer}. + * + * @author Anders Clausen + * @author Flaviu Muresan + * @since 3.1.0 + */ +public class TracingKafkaProducerBeanPostProcessor implements BeanPostProcessor { + + private final BeanFactory beanFactory; + + public TracingKafkaProducerBeanPostProcessor(BeanFactory beanFactory) { + this.beanFactory = beanFactory; + } + + @Override + public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { + if (bean instanceof Producer && !(bean instanceof TracingKafkaProducer)) { + return new TracingKafkaProducer<>((Producer) bean, beanFactory); + } + return bean; + } + +} diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingReactorKafkaAutoConfiguration.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingReactorKafkaAutoConfiguration.java new file mode 100644 index 000000000..49615622f --- /dev/null +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/TracingReactorKafkaAutoConfiguration.java @@ -0,0 +1,61 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.autoconfig.instrument.kafka; + +import reactor.kafka.receiver.KafkaReceiver; + +import org.springframework.beans.factory.BeanFactory; +import org.springframework.boot.autoconfigure.AutoConfigureAfter; +import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.autoconfig.brave.BraveAutoConfiguration; +import org.springframework.cloud.sleuth.instrument.kafka.TracingKafkaConsumerFactory; +import org.springframework.cloud.sleuth.instrument.kafka.TracingKafkaProducerFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +/** + * {@link org.springframework.boot.autoconfigure.EnableAutoConfiguration + * Auto-configuration} that registers instrumentation for Reactor Kafka. + * + * @author Anders Clausen + * @author Flaviu Muresan + * @since 3.1.0 + */ +@Configuration(proxyBeanMethods = false) +@ConditionalOnClass(KafkaReceiver.class) +@ConditionalOnBean(Tracer.class) +@AutoConfigureAfter(BraveAutoConfiguration.class) +@ConditionalOnProperty(value = "spring.sleuth.kafka.enabled", matchIfMissing = true) +public class TracingReactorKafkaAutoConfiguration { + + @Bean + @ConditionalOnMissingBean + TracingKafkaProducerFactory tracingKafkaProducerFactory(BeanFactory beanFactory) { + return new TracingKafkaProducerFactory(beanFactory); + } + + @Bean + @ConditionalOnMissingBean + TracingKafkaConsumerFactory tracingKafkaConsumerFactory(BeanFactory beanFactory) { + return new TracingKafkaConsumerFactory(beanFactory); + } + +} diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/messaging/TraceSpringMessagingAutoConfiguration.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/messaging/TraceSpringMessagingAutoConfiguration.java index a143bdb75..47f4a2d6f 100644 --- a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/messaging/TraceSpringMessagingAutoConfiguration.java +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/messaging/TraceSpringMessagingAutoConfiguration.java @@ -45,13 +45,13 @@ class TraceSpringMessagingAutoConfiguration { } @Bean - @ConditionalOnMissingBean + @ConditionalOnMissingBean(value = MessageHeaderAccessor.class, parameterizedContainer = Propagator.Setter.class) Propagator.Setter traceMessagePropagationSetter() { return new MessageHeaderPropagatorSetter(); } @Bean - @ConditionalOnMissingBean + @ConditionalOnMissingBean(value = MessageHeaderAccessor.class, parameterizedContainer = Propagator.Getter.class) Propagator.Getter traceMessagePropagationGetter() { return new MessageHeaderPropagatorGetter(); } diff --git a/spring-cloud-sleuth-autoconfigure/src/main/resources/META-INF/additional-spring-configuration-metadata.json b/spring-cloud-sleuth-autoconfigure/src/main/resources/META-INF/additional-spring-configuration-metadata.json index c94fec07a..9a894ba23 100644 --- a/spring-cloud-sleuth-autoconfigure/src/main/resources/META-INF/additional-spring-configuration-metadata.json +++ b/spring-cloud-sleuth-autoconfigure/src/main/resources/META-INF/additional-spring-configuration-metadata.json @@ -41,6 +41,12 @@ "description": "Enable tracing for WebSockets.", "defaultValue": true }, + { + "name": "spring.sleuth.kafka.enabled", + "type": "java.lang.Boolean", + "description": "Enable instrumenting of Apache Kafka clients.", + "defaultValue": true + }, { "name": "spring.sleuth.async.enabled", "type": "java.lang.Boolean", diff --git a/spring-cloud-sleuth-autoconfigure/src/main/resources/META-INF/spring.factories b/spring-cloud-sleuth-autoconfigure/src/main/resources/META-INF/spring.factories index d91d9d53e..7d7e98a50 100644 --- a/spring-cloud-sleuth-autoconfigure/src/main/resources/META-INF/spring.factories +++ b/spring-cloud-sleuth-autoconfigure/src/main/resources/META-INF/spring.factories @@ -1,5 +1,6 @@ # Auto Configuration org.springframework.boot.autoconfigure.EnableAutoConfiguration=\ +org.springframework.cloud.sleuth.autoconfig.instrument.kafka.TracingKafkaAutoConfiguration,\ org.springframework.cloud.sleuth.autoconfig.instrument.async.TraceAsyncAutoConfiguration,\ org.springframework.cloud.sleuth.autoconfig.instrument.async.TraceAsyncCustomAutoConfiguration,\ org.springframework.cloud.sleuth.autoconfig.instrument.async.TraceAsyncDefaultAutoConfiguration,\ diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTracingCallback.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTracingCallback.java index 4e8dfd44a..77da0902f 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTracingCallback.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTracingCallback.java @@ -31,7 +31,7 @@ import org.springframework.cloud.sleuth.Tracer; * * @author Anders Clausen * @author Flaviu Muresan - * @since 3.0.3 + * @since 3.1.0 */ public class KafkaTracingCallback implements Callback { diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumer.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumer.java index 9698f15cd..4e655899e 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumer.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumer.java @@ -37,8 +37,11 @@ import org.apache.kafka.common.MetricName; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; +import org.springframework.beans.factory.BeanFactory; import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.core.ParameterizedTypeReference; +import org.springframework.core.ResolvableType; /** * This decorates a Kafka {@link Consumer}. It creates and completes a @@ -47,21 +50,39 @@ import org.springframework.cloud.sleuth.propagation.Propagator; * * @author Anders Clausen * @author Flaviu Muresan - * @since 3.0.3 + * @since 3.1.0 */ public class TracingKafkaConsumer implements Consumer { + private final BeanFactory beanFactory; + private final Consumer delegate; - private final Propagator propagator; + private Propagator propagator; - private final Propagator.Getter> extractor; + private Propagator.Getter> extractor; - public TracingKafkaConsumer(Consumer consumer, Propagator propagator, - Propagator.Getter> getter) { + public TracingKafkaConsumer(Consumer consumer, BeanFactory beanFactory) { this.delegate = consumer; - this.propagator = propagator; - this.extractor = getter; + this.beanFactory = beanFactory; + } + + private Propagator propagator() { + if (this.propagator == null) { + this.propagator = this.beanFactory.getBean(Propagator.class); + } + return this.propagator; + } + + private Propagator.Getter> extractor() { + if (this.extractor == null) { + this.extractor = (Propagator.Getter>) beanFactory + .getBeanProvider(ResolvableType.forClassWithGenerics(Propagator.Getter.class, + ResolvableType.forType(new ParameterizedTypeReference>() { + }))) + .getIfAvailable(); + } + return this.extractor; } @Override @@ -109,7 +130,7 @@ public class TracingKafkaConsumer implements Consumer { public ConsumerRecords poll(long l) { ConsumerRecords consumerRecords = this.delegate.poll(l); for (ConsumerRecord consumerRecord : consumerRecords) { - KafkaTracingUtils.buildAndFinishSpan(consumerRecord, this.propagator, this.extractor); + KafkaTracingUtils.buildAndFinishSpan(consumerRecord, propagator(), extractor()); } return consumerRecords; } @@ -118,7 +139,7 @@ public class TracingKafkaConsumer implements Consumer { public ConsumerRecords poll(Duration duration) { ConsumerRecords consumerRecords = this.delegate.poll(duration); for (ConsumerRecord consumerRecord : consumerRecords) { - KafkaTracingUtils.buildAndFinishSpan(consumerRecord, this.propagator, this.extractor); + KafkaTracingUtils.buildAndFinishSpan(consumerRecord, propagator(), extractor()); } return consumerRecords; } diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumerFactory.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumerFactory.java new file mode 100644 index 000000000..342b669b0 --- /dev/null +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumerFactory.java @@ -0,0 +1,49 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.instrument.kafka; + +import org.apache.kafka.clients.consumer.Consumer; +import reactor.kafka.receiver.KafkaReceiver; +import reactor.kafka.receiver.ReceiverOptions; +import reactor.kafka.receiver.internals.ConsumerFactory; + +import org.springframework.beans.factory.BeanFactory; + +/** + * This decorates a Reactor Kafka {@link ConsumerFactory} to create decorated consumers of + * type {@link TracingKafkaConsumer}. This can be used by the {@link KafkaReceiver} + * factory methods to create instrumented receivers. + * + * @author Anders Clausen + * @author Flaviu Muresan + * @since 3.1.0 + */ +public class TracingKafkaConsumerFactory extends ConsumerFactory { + + private final BeanFactory beanFactory; + + public TracingKafkaConsumerFactory(BeanFactory beanFactory) { + super(); + this.beanFactory = beanFactory; + } + + @Override + public Consumer createConsumer(ReceiverOptions receiverOptions) { + return new TracingKafkaConsumer<>(super.createConsumer(receiverOptions), beanFactory); + } + +} diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducer.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducer.java index 51d778076..9a75e6748 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducer.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducer.java @@ -35,9 +35,12 @@ import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.ProducerFencedException; +import org.springframework.beans.factory.BeanFactory; import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.Tracer; import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.core.ParameterizedTypeReference; +import org.springframework.core.ResolvableType; /** * This decorates a Kafka {@link Producer} and creates a {@link Span.Kind#PRODUCER} span @@ -46,26 +49,50 @@ import org.springframework.cloud.sleuth.propagation.Propagator; * * @author Anders Clausen * @author Flaviu Muresan - * @since 3.0.3 + * @since 3.1.0 */ public class TracingKafkaProducer implements Producer { private static final Log log = LogFactory.getLog(TracingKafkaProducer.class); + private final BeanFactory beanFactory; + private final Producer delegate; - private final Tracer tracer; + private Tracer tracer; - private final Propagator propagator; + private Propagator propagator; - private final Propagator.Setter> injector; + private Propagator.Setter> injector; - public TracingKafkaProducer(Producer producer, Tracer tracer, Propagator propagator, - Propagator.Setter> setter) { + public TracingKafkaProducer(Producer producer, BeanFactory beanFactory) { this.delegate = producer; - this.tracer = tracer; - this.propagator = propagator; - this.injector = setter; + this.beanFactory = beanFactory; + } + + private Tracer tracer() { + if (this.tracer == null) { + this.tracer = this.beanFactory.getBean(Tracer.class); + } + return this.tracer; + } + + private Propagator propagator() { + if (this.propagator == null) { + this.propagator = this.beanFactory.getBean(Propagator.class); + } + return this.propagator; + } + + private Propagator.Setter> injector() { + if (this.injector == null) { + this.injector = (Propagator.Setter>) beanFactory + .getBeanProvider(ResolvableType.forClassWithGenerics(Propagator.Setter.class, + ResolvableType.forType(new ParameterizedTypeReference>() { + }))) + .getIfAvailable(); + } + return this.injector; } @Override @@ -107,15 +134,15 @@ public class TracingKafkaProducer implements Producer { @Override public Future send(ProducerRecord producerRecord, Callback callback) { - Span.Builder spanBuilder = tracer.spanBuilder().kind(Span.Kind.PRODUCER).name("kafka.produce") + Span.Builder spanBuilder = tracer().spanBuilder().kind(Span.Kind.PRODUCER).name("kafka.produce") .tag("kafka.topic", producerRecord.topic()); Span span = spanBuilder.start(); - this.propagator.inject(span.context(), producerRecord, this.injector); - try (Tracer.SpanInScope spanInScope = tracer.withSpan(span)) { + propagator().inject(span.context(), producerRecord, injector()); + try (Tracer.SpanInScope spanInScope = tracer().withSpan(span)) { if (log.isDebugEnabled()) { log.debug("Created producer span " + span); } - return this.delegate.send(producerRecord, new KafkaTracingCallback(callback, tracer, span)); + return this.delegate.send(producerRecord, new KafkaTracingCallback(callback, tracer(), span)); } } diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducerFactory.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducerFactory.java index 184df69b7..0f2200483 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducerFactory.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducerFactory.java @@ -21,8 +21,7 @@ import reactor.kafka.sender.KafkaSender; import reactor.kafka.sender.SenderOptions; import reactor.kafka.sender.internals.ProducerFactory; -import org.springframework.cloud.sleuth.Tracer; -import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.beans.factory.BeanFactory; /** * This decorates a Reactor Kafka {@link ProducerFactory} to create decorated producers of @@ -31,24 +30,20 @@ import org.springframework.cloud.sleuth.propagation.Propagator; * * @author Anders Clausen * @author Flaviu Muresan - * @since 3.0.3 + * @since 3.1.0 */ public class TracingKafkaProducerFactory extends ProducerFactory { - private final Tracer tracer; + private final BeanFactory beanFactory; - private final Propagator propagator; - - public TracingKafkaProducerFactory(Tracer tracer, Propagator propagator) { + public TracingKafkaProducerFactory(BeanFactory beanFactory) { super(); - this.tracer = tracer; - this.propagator = propagator; + this.beanFactory = beanFactory; } @Override public Producer createProducer(SenderOptions senderOptions) { - return new TracingKafkaProducer<>(super.createProducer(senderOptions), tracer, propagator, - new TracingKafkaPropagatorSetter()); + return new TracingKafkaProducer<>(super.createProducer(senderOptions), beanFactory); } } diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaPropagatorGetter.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaPropagatorGetter.java index e09988e5c..b5a96cebf 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaPropagatorGetter.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaPropagatorGetter.java @@ -30,7 +30,7 @@ import org.springframework.cloud.sleuth.propagation.Propagator; * * @author Anders Clausen * @author Flaviu Muresan - * @since 3.0.3 + * @since 3.1.0 */ public class TracingKafkaPropagatorGetter implements Propagator.Getter> { diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaPropagatorSetter.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaPropagatorSetter.java index afdb9b955..d327235ca 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaPropagatorSetter.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaPropagatorSetter.java @@ -26,7 +26,7 @@ import org.springframework.cloud.sleuth.propagation.Propagator; * * @author Anders Clausen * @author Flaviu Muresan - * @since 3.0.3 + * @since 3.1.0 */ public class TracingKafkaPropagatorSetter implements Propagator.Setter> { diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaReceiver.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaReceiver.java deleted file mode 100644 index f855b5595..000000000 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaReceiver.java +++ /dev/null @@ -1,112 +0,0 @@ -/* - * Copyright 2013-2021 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. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.sleuth.instrument.kafka; - -import java.util.function.Function; - -import org.apache.kafka.clients.consumer.Consumer; -import org.apache.kafka.clients.consumer.ConsumerRecord; -import reactor.core.publisher.Flux; -import reactor.core.publisher.Mono; -import reactor.kafka.receiver.KafkaReceiver; -import reactor.kafka.receiver.ReceiverRecord; -import reactor.kafka.sender.TransactionManager; - -import org.springframework.cloud.sleuth.Span; -import org.springframework.cloud.sleuth.propagation.Propagator; - -/** - * This decorates a reactive {@link KafkaReceiver} and creates and completes a - * {@link Span.Kind#CONSUMER} span for each record received. This span will be a child - * span of the one extracted from the record headers. - * - * @author Anders Clausen - * @author Flaviu Muresan - * @since 3.0.3 - */ -public class TracingKafkaReceiver implements KafkaReceiver { - - private final KafkaReceiver delegate; - - private final Propagator propagator; - - private final Propagator.Getter> extractor; - - public TracingKafkaReceiver(KafkaReceiver receiver, Propagator propagator, - Propagator.Getter> getter) { - this.delegate = receiver; - this.propagator = propagator; - this.extractor = getter; - } - - @Override - public Flux> receive(Integer integer) { - return buildAndFinishSpanOnNextReceiverRecord(this.delegate.receive(integer)); - } - - @Override - public Flux> receive() { - return buildAndFinishSpanOnNextReceiverRecord(this.delegate.receive()); - } - - @Override - public Flux>> receiveAutoAck(Integer integer) { - return this.delegate.receiveAutoAck(integer).map(this::buildAndFinishSpanOnNextConsumerRecord); - } - - @Override - public Flux>> receiveAutoAck() { - return this.delegate.receiveAutoAck().map(this::buildAndFinishSpanOnNextConsumerRecord); - } - - @Override - public Flux> receiveAtmostOnce(Integer integer) { - return this.buildAndFinishSpanOnNextConsumerRecord(this.delegate.receiveAtmostOnce(integer)); - } - - @Override - public Flux> receiveAtmostOnce() { - return this.buildAndFinishSpanOnNextConsumerRecord(this.delegate.receiveAtmostOnce()); - } - - @Override - public Flux>> receiveExactlyOnce(TransactionManager transactionManager) { - return this.delegate.receiveExactlyOnce(transactionManager).map(this::buildAndFinishSpanOnNextConsumerRecord); - } - - @Override - public Flux>> receiveExactlyOnce(TransactionManager transactionManager, Integer integer) { - return this.delegate.receiveExactlyOnce(transactionManager, integer) - .map(this::buildAndFinishSpanOnNextConsumerRecord); - } - - @Override - public Mono doOnConsumer(Function, ? extends T> function) { - return this.delegate.doOnConsumer(function); - } - - private Flux> buildAndFinishSpanOnNextConsumerRecord(Flux> flux) { - return flux.doOnNext(consumerRecord -> KafkaTracingUtils.buildAndFinishSpan(consumerRecord, this.propagator, - this.extractor)); - } - - private Flux> buildAndFinishSpanOnNextReceiverRecord(Flux> flux) { - return flux.doOnNext(consumerRecord -> KafkaTracingUtils.buildAndFinishSpan(consumerRecord, this.propagator, - this.extractor)); - } - -} diff --git a/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumerTest.java b/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumerTest.java index c04c22d72..dcac1f099 100644 --- a/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumerTest.java +++ b/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaConsumerTest.java @@ -35,6 +35,8 @@ import org.mockito.Mock; import org.mockito.Mockito; import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.support.StaticListableBeanFactory; import org.springframework.cloud.sleuth.propagation.Propagator; import static org.mockito.ArgumentMatchers.eq; @@ -48,6 +50,9 @@ public class TracingKafkaConsumerTest { @Mock(answer = Answers.RETURNS_DEEP_STUBS) Propagator propagator; + @Mock + Propagator.Getter> extractor; + @Test void should_delegate_poll_calls() { Duration pollTimeout = Duration.of(5, ChronoUnit.SECONDS); @@ -57,11 +62,18 @@ public class TracingKafkaConsumerTest { ConsumerRecords records = new ConsumerRecords<>(map); BDDMockito.given(kafkaConsumer.poll(pollTimeout)).willReturn(records); TracingKafkaConsumer tracingKafkaConsumer = new TracingKafkaConsumer<>(kafkaConsumer, - propagator, new TracingKafkaPropagatorGetter()); + beanFactory()); tracingKafkaConsumer.poll(pollTimeout); Mockito.verify(kafkaConsumer).poll(eq(pollTimeout)); } + private BeanFactory beanFactory() { + StaticListableBeanFactory beanFactory = new StaticListableBeanFactory(); + beanFactory.addBean("propagator", this.propagator); + beanFactory.addBean("extractor", this.extractor); + return beanFactory; + } + } diff --git a/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducerTest.java b/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducerTest.java index dcb3bab8b..568968727 100644 --- a/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducerTest.java +++ b/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaProducerTest.java @@ -28,6 +28,8 @@ import org.mockito.Mock; import org.mockito.Mockito; import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.support.StaticListableBeanFactory; import org.springframework.cloud.sleuth.Tracer; import org.springframework.cloud.sleuth.propagation.Propagator; import org.springframework.test.util.ReflectionTestUtils; @@ -52,8 +54,8 @@ public class TracingKafkaProducerTest { ProducerRecord testRecord = new ProducerRecord<>("test", "test"); Callback callback = (record, ex) -> { }; - TracingKafkaProducer tracingKafkaProducer = new TracingKafkaProducer<>(kafkaProducer, tracer, - propagator, new TracingKafkaPropagatorSetter()); + TracingKafkaProducer tracingKafkaProducer = new TracingKafkaProducer<>(kafkaProducer, + beanFactory()); tracingKafkaProducer.send(testRecord, callback); @@ -65,8 +67,8 @@ public class TracingKafkaProducerTest { ProducerRecord testRecord = new ProducerRecord<>("test", "test"); Callback callback = (record, ex) -> { }; - TracingKafkaProducer tracingKafkaProducer = new TracingKafkaProducer<>(kafkaProducer, tracer, - propagator, new TracingKafkaPropagatorSetter()); + TracingKafkaProducer tracingKafkaProducer = new TracingKafkaProducer<>(kafkaProducer, + beanFactory()); tracingKafkaProducer.send(testRecord, callback); @@ -76,4 +78,11 @@ public class TracingKafkaProducerTest { BDDAssertions.then(ReflectionTestUtils.getField(callbackArgument.getValue(), "callback")).isEqualTo(callback); } + private BeanFactory beanFactory() { + StaticListableBeanFactory beanFactory = new StaticListableBeanFactory(); + beanFactory.addBean("tracer", this.tracer); + beanFactory.addBean("propagator", this.propagator); + return beanFactory; + } + } diff --git a/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaReceiverTest.java b/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaReceiverTest.java deleted file mode 100644 index 6707aeee6..000000000 --- a/spring-cloud-sleuth-instrumentation/src/test/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaReceiverTest.java +++ /dev/null @@ -1,61 +0,0 @@ -/* - * Copyright 2013-2021 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. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.sleuth.instrument.kafka; - -import java.util.function.Predicate; - -import org.apache.kafka.clients.consumer.ConsumerRecord; -import org.junit.jupiter.api.Test; -import org.junit.jupiter.api.extension.ExtendWith; -import org.mockito.Answers; -import org.mockito.BDDMockito; -import org.mockito.Mock; -import org.mockito.Mockito; -import org.mockito.junit.jupiter.MockitoExtension; -import reactor.core.publisher.Flux; -import reactor.kafka.receiver.KafkaReceiver; -import reactor.kafka.receiver.ReceiverOffset; -import reactor.kafka.receiver.ReceiverRecord; -import reactor.test.StepVerifier; - -import org.springframework.cloud.sleuth.propagation.Propagator; - -@ExtendWith(MockitoExtension.class) -public class TracingKafkaReceiverTest { - - @Mock - KafkaReceiver kafkaReceiver; - - @Mock(answer = Answers.RETURNS_DEEP_STUBS) - Propagator propagator; - - @Test - void should_delegate_receive_calls() { - ReceiverOffset receiverOffset = BDDMockito.mock(ReceiverOffset.class); - ConsumerRecord record = new ConsumerRecord<>("topic", 0, 1, "test-key", "test-value"); - ReceiverRecord receiverRecord = new ReceiverRecord<>(record, receiverOffset); - BDDMockito.given(kafkaReceiver.receive()).willReturn(Flux.just(receiverRecord)); - TracingKafkaReceiver tracingKafkaReceiver = new TracingKafkaReceiver<>(kafkaReceiver, - propagator, new TracingKafkaPropagatorGetter()); - - StepVerifier.create(tracingKafkaReceiver.receive()).expectNextMatches(Predicate.isEqual(receiverRecord)) - .verifyComplete(); - - Mockito.verify(kafkaReceiver).receive(); - } - -} diff --git a/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/pom.xml b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/pom.xml index 6faff3f17..5692dd97b 100644 --- a/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/pom.xml +++ b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/pom.xml @@ -70,7 +70,10 @@ io.projectreactor.kafka reactor-kafka - true + + + io.projectreactor + reactor-test org.testcontainers diff --git a/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaConsumerTest.java b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaConsumerTest.java new file mode 100644 index 000000000..fbe5cca26 --- /dev/null +++ b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaConsumerTest.java @@ -0,0 +1,65 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.brave.instrument.kafka; + +import java.time.Duration; + +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.assertj.core.api.BDDAssertions; +import org.junit.jupiter.api.Test; + +import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.brave.BraveTestTracing; +import org.springframework.cloud.sleuth.exporter.FinishedSpan; +import org.springframework.cloud.sleuth.instrument.kafka.KafkaTestUtils; +import org.springframework.cloud.sleuth.test.TestTracingAware; + +import static org.awaitility.Awaitility.await; + +public class KafkaConsumerTest extends org.springframework.cloud.sleuth.instrument.kafka.KafkaConsumerTest { + + BraveTestTracing testTracing; + + @Override + public TestTracingAware tracerTest() { + if (this.testTracing == null) { + this.testTracing = new BraveTestTracing(); + } + return this.testTracing; + } + + @Test + public void should_consider_native_headers() { + KafkaProducer kafkaProducer = KafkaTestUtils + .buildTestKafkaProducer(kafkaContainer.getBootstrapServers()); + ProducerRecord producerRecord = new ProducerRecord<>(testTopic, "test", "test"); + producerRecord.headers().add("b3", "000000000000000a-000000000000000b-1-000000000000000a".getBytes()); + kafkaProducer.send(producerRecord); + kafkaProducer.close(); + + await().atMost(Duration.ofSeconds(5)).until(() -> receivedCounter.intValue() == 1); + + BDDAssertions.then(this.tracer.currentSpan()).isNull(); + BDDAssertions.then(this.spans).hasSize(1); + FinishedSpan span = this.spans.get(0); + BDDAssertions.then(span.getKind()).isEqualTo(Span.Kind.CONSUMER); + BDDAssertions.then(span.getTraceId()).isEqualTo("000000000000000a"); + BDDAssertions.then(span.getParentId()).isEqualTo("000000000000000b"); + } + +} diff --git a/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaProducerTest.java b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaProducerTest.java index b19b6efc3..78f9e7a6e 100644 --- a/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaProducerTest.java +++ b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaProducerTest.java @@ -16,6 +16,15 @@ package org.springframework.cloud.sleuth.brave.instrument.kafka; +import java.util.Optional; +import java.util.concurrent.TimeUnit; + +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.header.Header; +import org.assertj.core.api.BDDAssertions; +import org.junit.jupiter.api.Test; + import org.springframework.cloud.sleuth.brave.BraveTestTracing; import org.springframework.cloud.sleuth.test.TestTracingAware; @@ -31,4 +40,23 @@ public class KafkaProducerTest extends org.springframework.cloud.sleuth.instrume return this.testTracing; } + @Test + public void should_inject_native_headers() throws InterruptedException { + ProducerRecord producerRecord = new ProducerRecord<>(testTopic, "test", "test"); + startKafkaConsumer(); + + this.kafkaProducer.send(producerRecord); + ConsumerRecord consumerRecord = consumerRecords.poll(5, TimeUnit.SECONDS); + + BDDAssertions.then(consumerRecord).isNotNull(); + BDDAssertions.then(getHeaderValueOrNull(consumerRecord, "X-B3-TraceId")).isNotNull(); + BDDAssertions.then(getHeaderValueOrNull(consumerRecord, "X-B3-SpanId")).isNotNull(); + BDDAssertions.then(getHeaderValueOrNull(consumerRecord, "X-B3-Sampled")).isNotNull(); + } + + private static String getHeaderValueOrNull(ConsumerRecord consumerRecord, String header) { + return Optional.ofNullable(consumerRecord).map(ConsumerRecord::headers) + .map(headers -> headers.lastHeader(header)).map(Header::value).map(String::new).orElse(null); + } + } diff --git a/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaReceiverTest.java b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaReceiverTest.java new file mode 100644 index 000000000..5a3f14d75 --- /dev/null +++ b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaReceiverTest.java @@ -0,0 +1,64 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.brave.instrument.kafka; + +import java.time.Duration; + +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.assertj.core.api.BDDAssertions; +import org.junit.jupiter.api.Test; + +import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.brave.BraveTestTracing; +import org.springframework.cloud.sleuth.exporter.FinishedSpan; +import org.springframework.cloud.sleuth.instrument.kafka.KafkaTestUtils; +import org.springframework.cloud.sleuth.test.TestTracingAware; + +import static org.awaitility.Awaitility.await; + +public class KafkaReceiverTest extends org.springframework.cloud.sleuth.instrument.kafka.KafkaReceiverTest { + + BraveTestTracing testTracing; + + @Override + public TestTracingAware tracerTest() { + if (this.testTracing == null) { + this.testTracing = new BraveTestTracing(); + } + return this.testTracing; + } + + @Test + public void should_consider_native_headers() { + KafkaProducer kafkaProducer = KafkaTestUtils + .buildTestKafkaProducer(kafkaContainer.getBootstrapServers()); + ProducerRecord producerRecord = new ProducerRecord<>(testTopic, "test", "test"); + producerRecord.headers().add("b3", "000000000000000a-000000000000000b-1-000000000000000a".getBytes()); + kafkaProducer.send(producerRecord); + + await().atMost(Duration.ofSeconds(5)).until(() -> receivedCounter.intValue() == 1); + + BDDAssertions.then(this.tracer.currentSpan()).isNull(); + BDDAssertions.then(this.spans).hasSize(1); + FinishedSpan span = this.spans.get(0); + BDDAssertions.then(span.getKind()).isEqualTo(Span.Kind.CONSUMER); + BDDAssertions.then(span.getTraceId()).isEqualTo("000000000000000a"); + BDDAssertions.then(span.getParentId()).isEqualTo("000000000000000b"); + } + +} diff --git a/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaSenderTest.java b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaSenderTest.java new file mode 100644 index 000000000..fb88a5b99 --- /dev/null +++ b/tests/brave/spring-cloud-sleuth-instrumentation-kafka-tests/src/test/java/org/springframework/cloud/sleuth/brave/instrument/kafka/KafkaSenderTest.java @@ -0,0 +1,69 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.brave.instrument.kafka; + +import java.util.Optional; +import java.util.concurrent.TimeUnit; + +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.header.Header; +import org.assertj.core.api.BDDAssertions; +import org.junit.jupiter.api.Test; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.kafka.sender.SenderRecord; +import reactor.kafka.sender.SenderResult; +import reactor.test.StepVerifier; + +import org.springframework.cloud.sleuth.brave.BraveTestTracing; +import org.springframework.cloud.sleuth.test.TestTracingAware; + +public class KafkaSenderTest extends org.springframework.cloud.sleuth.instrument.kafka.KafkaSenderTest { + + BraveTestTracing testTracing; + + @Override + public TestTracingAware tracerTest() { + if (this.testTracing == null) { + this.testTracing = new BraveTestTracing(); + } + return this.testTracing; + } + + @Test + public void should_inject_native_headers() throws InterruptedException { + ProducerRecord producerRecord = new ProducerRecord<>(testTopic, "test", "test"); + startKafkaConsumer(); + + Flux> senderResultFlux = this.kafkaSender + .send(Mono.just(SenderRecord.create(producerRecord, null))); + StepVerifier.create(senderResultFlux).expectNextCount(1).verifyComplete(); + ConsumerRecord consumerRecord = consumerRecords.poll(5, TimeUnit.SECONDS); + + BDDAssertions.then(consumerRecord).isNotNull(); + BDDAssertions.then(getHeaderValueOrNull(consumerRecord, "X-B3-TraceId")).isNotNull(); + BDDAssertions.then(getHeaderValueOrNull(consumerRecord, "X-B3-SpanId")).isNotNull(); + BDDAssertions.then(getHeaderValueOrNull(consumerRecord, "X-B3-Sampled")).isNotNull(); + } + + private static String getHeaderValueOrNull(ConsumerRecord consumerRecord, String header) { + return Optional.ofNullable(consumerRecord).map(ConsumerRecord::headers) + .map(headers -> headers.lastHeader(header)).map(Header::value).map(String::new).orElse(null); + } + +} diff --git a/tests/common/pom.xml b/tests/common/pom.xml index a4f029dc5..3f6a1176d 100644 --- a/tests/common/pom.xml +++ b/tests/common/pom.xml @@ -149,6 +149,11 @@ brave-tests true + + io.projectreactor + reactor-test + true + io.projectreactor.kafka reactor-kafka diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaConsumerTest.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaConsumerTest.java new file mode 100644 index 000000000..cd84d04b3 --- /dev/null +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaConsumerTest.java @@ -0,0 +1,159 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.instrument.kafka; + +import java.time.Duration; +import java.util.HashMap; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.regex.Pattern; + +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.assertj.core.api.BDDAssertions; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Answers; +import org.mockito.BDDMockito; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.testcontainers.containers.KafkaContainer; +import org.testcontainers.containers.wait.strategy.Wait; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +import org.springframework.beans.factory.BeanFactory; +import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.exporter.FinishedSpan; +import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.cloud.sleuth.test.TestSpanHandler; +import org.springframework.cloud.sleuth.test.TestTracingAwareSupplier; +import org.springframework.core.ParameterizedTypeReference; +import org.springframework.core.ResolvableType; + +import static org.awaitility.Awaitility.await; + +@Testcontainers +@ExtendWith(MockitoExtension.class) +public abstract class KafkaConsumerTest implements TestTracingAwareSupplier { + + protected String testTopic; + + protected Tracer tracer = tracerTest().tracing().tracer(); + + protected Propagator propagator = tracerTest().tracing().propagator(); + + protected TestSpanHandler spans = tracerTest().handler(); + + protected TracingKafkaConsumer kafkaConsumer; + + private final AtomicBoolean consumerRun = new AtomicBoolean(); + + protected final AtomicInteger receivedCounter = new AtomicInteger(0); + + @Mock(answer = Answers.RETURNS_DEEP_STUBS) + BeanFactory beanFactory; + + @Container + protected static final KafkaContainer kafkaContainer = new KafkaContainer( + DockerImageName.parse("confluentinc/cp-kafka:6.1.1")).withExposedPorts(9093) + .waitingFor(Wait.forListeningPort()); + + @BeforeAll + static void setupAll() { + kafkaContainer.start(); + } + + @AfterAll + static void destroyAll() { + kafkaContainer.stop(); + } + + @BeforeEach + void setup() { + BDDMockito.given(this.beanFactory.getBean(Propagator.class)).willReturn(this.propagator); + BDDMockito.given(this.beanFactory.getBeanProvider(ResolvableType.forClassWithGenerics(Propagator.Getter.class, + ResolvableType.forType(new ParameterizedTypeReference>() { + }))).getIfAvailable()).willReturn(new TracingKafkaPropagatorGetter()); + testTopic = UUID.randomUUID().toString(); + Map consumerProperties = new HashMap<>(); + consumerProperties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers()); + consumerProperties.put(ConsumerConfig.GROUP_ID_CONFIG, "test-consumer-group"); + consumerProperties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProperties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + kafkaConsumer = new TracingKafkaConsumer<>(new KafkaConsumer<>(consumerProperties), beanFactory); + consumerRun.set(true); + Executors.newSingleThreadExecutor().execute(() -> doStartKafkaConsumer(receivedCounter)); + } + + @AfterEach + void destroy() { + consumerRun.set(false); + } + + @Test + public void should_create_and_finish_consumer_span() { + KafkaProducer kafkaProducer = KafkaTestUtils + .buildTestKafkaProducer(kafkaContainer.getBootstrapServers()); + ProducerRecord producerRecord = new ProducerRecord<>(testTopic, "test", "test"); + kafkaProducer.send(producerRecord); + kafkaProducer.close(); + + await().atMost(Duration.ofSeconds(5)).until(() -> receivedCounter.intValue() == 1); + + BDDAssertions.then(this.tracer.currentSpan()).isNull(); + BDDAssertions.then(this.spans).hasSize(1); + FinishedSpan span = this.spans.get(0); + BDDAssertions.then(span.getKind()).isEqualTo(Span.Kind.CONSUMER); + BDDAssertions.then(span.getTags()).isNotEmpty(); + BDDAssertions.then(span.getTags().get("kafka.topic")).isEqualTo(testTopic); + BDDAssertions.then(span.getTags().get("kafka.offset")).isEqualTo("0"); + BDDAssertions.then(span.getTags().get("kafka.partition")).isEqualTo("0"); + } + + private void doStartKafkaConsumer(AtomicInteger receivedCounter) { + this.kafkaConsumer.subscribe(Pattern.compile(this.testTopic)); + while (this.consumerRun.get()) { + ConsumerRecords records = this.kafkaConsumer.poll(Duration.ofSeconds(1)); + for (ConsumerRecord record : records) { + receivedCounter.incrementAndGet(); + } + } + this.kafkaConsumer.close(); + } + + @Override + public void cleanUpTracing() { + this.spans.clear(); + } + +} diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaProducerTest.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaProducerTest.java index c541eac73..f399eb190 100644 --- a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaProducerTest.java +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaProducerTest.java @@ -19,34 +19,56 @@ package org.springframework.cloud.sleuth.instrument.kafka; import java.time.Duration; import java.util.HashMap; import java.util.Map; +import java.util.UUID; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.Executors; +import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.regex.Pattern; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.producer.Callback; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import org.assertj.core.api.BDDAssertions; +import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Answers; +import org.mockito.BDDMockito; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; import org.testcontainers.containers.KafkaContainer; import org.testcontainers.containers.wait.strategy.Wait; import org.testcontainers.junit.jupiter.Container; import org.testcontainers.junit.jupiter.Testcontainers; import org.testcontainers.utility.DockerImageName; +import org.springframework.beans.factory.BeanFactory; import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.exporter.FinishedSpan; import org.springframework.cloud.sleuth.propagation.Propagator; import org.springframework.cloud.sleuth.test.TestSpanHandler; import org.springframework.cloud.sleuth.test.TestTracingAwareSupplier; +import org.springframework.core.ParameterizedTypeReference; +import org.springframework.core.ResolvableType; import static org.awaitility.Awaitility.await; @Testcontainers +@ExtendWith(MockitoExtension.class) public abstract class KafkaProducerTest implements TestTracingAwareSupplier { + protected String testTopic; + protected Tracer tracer = tracerTest().tracing().tracer(); protected Propagator propagator = tracerTest().tracing().propagator(); @@ -55,39 +77,82 @@ public abstract class KafkaProducerTest implements TestTracingAwareSupplier { protected TracingKafkaProducer kafkaProducer; + private final AtomicBoolean consumerRun = new AtomicBoolean(); + + protected final BlockingQueue> consumerRecords = new LinkedBlockingQueue<>(); + + @Mock(answer = Answers.RETURNS_DEEP_STUBS) + BeanFactory beanFactory; + @Container - protected final KafkaContainer kafkaContainer = new KafkaContainer( - DockerImageName.parse("confluentinc/cp-kafka:5.2.1")).withExposedPorts(9093) + protected static final KafkaContainer kafkaContainer = new KafkaContainer( + DockerImageName.parse("confluentinc/cp-kafka:6.1.1")).withExposedPorts(9093) .waitingFor(Wait.forListeningPort()); + @BeforeAll + static void setupAll() { + kafkaContainer.start(); + } + + @AfterAll + static void destroyAll() { + kafkaContainer.stop(); + } + @BeforeEach void setup() { - kafkaContainer.start(); - Map properties = new HashMap<>(); - properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers()); - properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); - properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); - kafkaProducer = new TracingKafkaProducer<>(new KafkaProducer<>(properties), tracer, propagator, - new TracingKafkaPropagatorSetter()); + BDDMockito.given(this.beanFactory.getBean(Tracer.class)).willReturn(this.tracer); + BDDMockito.given(this.beanFactory.getBean(Propagator.class)).willReturn(this.propagator); + BDDMockito.given(this.beanFactory.getBeanProvider(ResolvableType.forClassWithGenerics(Propagator.Setter.class, + ResolvableType.forType(new ParameterizedTypeReference>() { + }))).getIfAvailable()).willReturn(new TracingKafkaPropagatorSetter()); + testTopic = UUID.randomUUID().toString(); + Map producerProperties = new HashMap<>(); + producerProperties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers()); + producerProperties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + producerProperties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + kafkaProducer = new TracingKafkaProducer<>(new KafkaProducer<>(producerProperties), beanFactory); + consumerRun.set(true); + consumerRecords.clear(); } @AfterEach void destroy() { - kafkaContainer.stop(); + this.kafkaProducer.close(); + consumerRun.set(false); } @Test public void should_create_and_finish_producer_span() { AtomicBoolean acknowledged = new AtomicBoolean(false); Callback callback = (metadata, ex) -> acknowledged.set(true); - ProducerRecord producerRecord = new ProducerRecord<>("spring-cloud-sleuth-otel-topic", "test", - "test"); + ProducerRecord producerRecord = new ProducerRecord<>(testTopic, "test", "test"); + this.kafkaProducer.send(producerRecord, callback); await().atMost(Duration.ofSeconds(5)).until(acknowledged::get); BDDAssertions.then(this.tracer.currentSpan()).isNull(); - BDDAssertions.then(this.spans).isNotEmpty(); - BDDAssertions.then(this.spans.get(0).getKind()).isEqualTo(Span.Kind.PRODUCER); + BDDAssertions.then(this.spans).hasSize(1); + FinishedSpan span = this.spans.get(0); + BDDAssertions.then(span.getKind()).isEqualTo(Span.Kind.PRODUCER); + BDDAssertions.then(span.getTags().get("kafka.topic")).isEqualTo(testTopic); + } + + protected void startKafkaConsumer() { + Executors.newSingleThreadExecutor().execute(this::doStartKafkaConsumer); + } + + private void doStartKafkaConsumer() { + KafkaConsumer kafkaConsumer = KafkaTestUtils + .buildTestKafkaConsumer(kafkaContainer.getBootstrapServers()); + kafkaConsumer.subscribe(Pattern.compile(testTopic)); + while (consumerRun.get()) { + ConsumerRecords records = kafkaConsumer.poll(Duration.ofSeconds(1)); + for (ConsumerRecord record : records) { + consumerRecords.offer(record); + } + } + kafkaConsumer.close(); } @Override diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaReceiverTest.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaReceiverTest.java new file mode 100644 index 000000000..e33f3eb7c --- /dev/null +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaReceiverTest.java @@ -0,0 +1,150 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.instrument.kafka; + +import java.time.Duration; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.assertj.core.api.BDDAssertions; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Answers; +import org.mockito.BDDMockito; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.testcontainers.containers.KafkaContainer; +import org.testcontainers.containers.wait.strategy.Wait; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; +import reactor.core.Disposable; +import reactor.core.scheduler.Schedulers; +import reactor.kafka.receiver.KafkaReceiver; +import reactor.kafka.receiver.ReceiverOptions; + +import org.springframework.beans.factory.BeanFactory; +import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.exporter.FinishedSpan; +import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.cloud.sleuth.test.TestSpanHandler; +import org.springframework.cloud.sleuth.test.TestTracingAwareSupplier; +import org.springframework.core.ParameterizedTypeReference; +import org.springframework.core.ResolvableType; + +import static org.awaitility.Awaitility.await; + +@Testcontainers +@ExtendWith(MockitoExtension.class) +public abstract class KafkaReceiverTest implements TestTracingAwareSupplier { + + protected String testTopic; + + protected Tracer tracer = tracerTest().tracing().tracer(); + + protected Propagator propagator = tracerTest().tracing().propagator(); + + protected TestSpanHandler spans = tracerTest().handler(); + + private Disposable consumerSubscription; + + protected final AtomicInteger receivedCounter = new AtomicInteger(0); + + @Mock(answer = Answers.RETURNS_DEEP_STUBS) + BeanFactory beanFactory; + + @Container + protected static final KafkaContainer kafkaContainer = new KafkaContainer( + DockerImageName.parse("confluentinc/cp-kafka:6.1.1")).withExposedPorts(9093) + .waitingFor(Wait.forListeningPort()); + + @BeforeAll + static void setupAll() { + kafkaContainer.start(); + } + + @AfterAll + static void destroyAll() { + kafkaContainer.stop(); + } + + @BeforeEach + void setup() { + BDDMockito.given(this.beanFactory.getBean(Propagator.class)).willReturn(this.propagator); + BDDMockito.given(this.beanFactory.getBeanProvider(ResolvableType.forClassWithGenerics(Propagator.Getter.class, + ResolvableType.forType(new ParameterizedTypeReference>() { + }))).getIfAvailable()).willReturn(new TracingKafkaPropagatorGetter()); + testTopic = UUID.randomUUID().toString(); + Map consumerProperties = new HashMap<>(); + consumerProperties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers()); + consumerProperties.put(ConsumerConfig.GROUP_ID_CONFIG, "test-consumer-group"); + consumerProperties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProperties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + ReceiverOptions options = ReceiverOptions.create(consumerProperties); + options = options.withKeyDeserializer(new StringDeserializer()).withValueDeserializer(new StringDeserializer()) + .subscription(Collections.singletonList(testTopic)); + KafkaReceiver kafkaReceiver = KafkaReceiver.create(new TracingKafkaConsumerFactory(beanFactory), + options); + this.consumerSubscription = kafkaReceiver.receive().subscribeOn(Schedulers.single()) + .subscribe(record -> receivedCounter.incrementAndGet()); + this.receivedCounter.set(0); + } + + @AfterEach + void destroy() { + this.consumerSubscription.dispose(); + } + + @Test + public void should_create_and_finish_consumer_span() { + KafkaProducer kafkaProducer = KafkaTestUtils + .buildTestKafkaProducer(kafkaContainer.getBootstrapServers()); + ProducerRecord producerRecord = new ProducerRecord<>(testTopic, "test", "test"); + kafkaProducer.send(producerRecord); + + await().atMost(Duration.ofSeconds(5)).until(() -> receivedCounter.intValue() == 1); + + BDDAssertions.then(this.tracer.currentSpan()).isNull(); + BDDAssertions.then(this.spans).hasSize(1); + FinishedSpan span = this.spans.get(0); + BDDAssertions.then(span.getKind()).isEqualTo(Span.Kind.CONSUMER); + BDDAssertions.then(span.getTags()).isNotEmpty(); + BDDAssertions.then(span.getTags().get("kafka.topic")).isEqualTo(testTopic); + BDDAssertions.then(span.getTags().get("kafka.offset")).isEqualTo("0"); + BDDAssertions.then(span.getTags().get("kafka.partition")).isEqualTo("0"); + } + + @Override + public void cleanUpTracing() { + this.spans.clear(); + } + +} diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaSenderTest.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaSenderTest.java new file mode 100644 index 000000000..67db8e9cc --- /dev/null +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaSenderTest.java @@ -0,0 +1,169 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.instrument.kafka; + +import java.time.Duration; +import java.util.HashMap; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.Executors; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.regex.Pattern; + +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.StringSerializer; +import org.assertj.core.api.BDDAssertions; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Answers; +import org.mockito.BDDMockito; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.testcontainers.containers.KafkaContainer; +import org.testcontainers.containers.wait.strategy.Wait; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.kafka.sender.KafkaSender; +import reactor.kafka.sender.SenderOptions; +import reactor.kafka.sender.SenderRecord; +import reactor.kafka.sender.SenderResult; +import reactor.test.StepVerifier; + +import org.springframework.beans.factory.BeanFactory; +import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.exporter.FinishedSpan; +import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.cloud.sleuth.test.TestSpanHandler; +import org.springframework.cloud.sleuth.test.TestTracingAwareSupplier; +import org.springframework.core.ParameterizedTypeReference; +import org.springframework.core.ResolvableType; + +@Testcontainers +@ExtendWith(MockitoExtension.class) +public abstract class KafkaSenderTest implements TestTracingAwareSupplier { + + protected String testTopic; + + protected Tracer tracer = tracerTest().tracing().tracer(); + + protected Propagator propagator = tracerTest().tracing().propagator(); + + protected TestSpanHandler spans = tracerTest().handler(); + + protected KafkaSender kafkaSender; + + private final AtomicBoolean consumerRun = new AtomicBoolean(); + + protected final BlockingQueue> consumerRecords = new LinkedBlockingQueue<>(); + + @Mock(answer = Answers.RETURNS_DEEP_STUBS) + BeanFactory beanFactory; + + @Container + protected static final KafkaContainer kafkaContainer = new KafkaContainer( + DockerImageName.parse("confluentinc/cp-kafka:6.1.1")).withExposedPorts(9093) + .waitingFor(Wait.forListeningPort()); + + @BeforeAll + static void setupAll() { + kafkaContainer.start(); + } + + @AfterAll + static void destroyAll() { + kafkaContainer.stop(); + } + + @BeforeEach + void setup() { + BDDMockito.given(this.beanFactory.getBean(Tracer.class)).willReturn(this.tracer); + BDDMockito.given(this.beanFactory.getBean(Propagator.class)).willReturn(this.propagator); + BDDMockito.given(this.beanFactory.getBeanProvider(ResolvableType.forClassWithGenerics(Propagator.Setter.class, + ResolvableType.forType(new ParameterizedTypeReference>() { + }))).getIfAvailable()).willReturn(new TracingKafkaPropagatorSetter()); + testTopic = UUID.randomUUID().toString(); + Map producerProperties = new HashMap<>(); + producerProperties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers()); + producerProperties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + producerProperties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + this.kafkaSender = KafkaSender.create(new TracingKafkaProducerFactory(beanFactory), + SenderOptions.create(producerProperties)); + consumerRun.set(true); + consumerRecords.clear(); + } + + @AfterEach + void destroy() { + consumerRun.set(false); + } + + @Test + public void should_create_and_finish_producer_span() throws InterruptedException { + ProducerRecord producerRecord = new ProducerRecord<>(testTopic, "test", "test"); + startKafkaConsumer(); + + Flux> senderResultFlux = this.kafkaSender + .send(Mono.just(SenderRecord.create(producerRecord, null))); + StepVerifier.create(senderResultFlux).expectNextCount(1).verifyComplete(); + consumerRecords.poll(5, TimeUnit.SECONDS); + + BDDAssertions.then(this.tracer.currentSpan()).isNull(); + BDDAssertions.then(this.spans).hasSize(1); + FinishedSpan span = this.spans.get(0); + BDDAssertions.then(span.getKind()).isEqualTo(Span.Kind.PRODUCER); + BDDAssertions.then(span.getTags().get("kafka.topic")).isEqualTo(testTopic); + + } + + protected void startKafkaConsumer() { + Executors.newSingleThreadExecutor().execute(this::doStartKafkaConsumer); + } + + private void doStartKafkaConsumer() { + KafkaConsumer kafkaConsumer = KafkaTestUtils + .buildTestKafkaConsumer(kafkaContainer.getBootstrapServers()); + kafkaConsumer.subscribe(Pattern.compile(this.testTopic)); + while (this.consumerRun.get()) { + ConsumerRecords records = kafkaConsumer.poll(Duration.ofSeconds(1)); + for (ConsumerRecord record : records) { + this.consumerRecords.offer(record); + } + } + kafkaConsumer.close(); + } + + @Override + public void cleanUpTracing() { + this.spans.clear(); + } + +} diff --git a/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTestUtils.java b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTestUtils.java new file mode 100644 index 000000000..4d58a2306 --- /dev/null +++ b/tests/common/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTestUtils.java @@ -0,0 +1,52 @@ +/* + * Copyright 2013-2021 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.instrument.kafka; + +import java.util.HashMap; +import java.util.Map; + +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.apache.kafka.common.serialization.StringSerializer; + +public final class KafkaTestUtils { + + private KafkaTestUtils() { + } + + public static KafkaProducer buildTestKafkaProducer(String bootstrapServers) { + Map producerProperties = new HashMap<>(); + producerProperties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); + producerProperties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + producerProperties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + return new KafkaProducer<>(producerProperties); + } + + public static KafkaConsumer buildTestKafkaConsumer(String bootstrapServers) { + Map consumerProperties = new HashMap<>(); + consumerProperties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); + consumerProperties.put(ConsumerConfig.GROUP_ID_CONFIG, "test-consumer-group"); + consumerProperties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProperties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + return new KafkaConsumer<>(consumerProperties); + } + +}