From 73f9ec61f61596da2b67e893d096322e2247b8c0 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Fri, 9 Mar 2018 10:16:26 +0100 Subject: [PATCH] Added spring-kafka support; fixes gh-896 --- .../main/asciidoc/spring-cloud-sleuth.adoc | 12 ++++- spring-cloud-sleuth-core/pom.xml | 9 ++++ .../TraceMessagingAutoConfiguration.java | 54 +++++++++++++++++++ .../TraceMessagingAutoConfigurationTests.java | 41 ++++++++++++++ 4 files changed, 115 insertions(+), 1 deletion(-) diff --git a/docs/src/main/asciidoc/spring-cloud-sleuth.adoc b/docs/src/main/asciidoc/spring-cloud-sleuth.adoc index 2021a92eb..74ad78a04 100644 --- a/docs/src/main/asciidoc/spring-cloud-sleuth.adoc +++ b/docs/src/main/asciidoc/spring-cloud-sleuth.adoc @@ -1125,6 +1125,8 @@ include::../../../../spring-cloud-sleuth-core/src/test/java/org/springframework/ === Messaging +==== Spring Integration and Spring Cloud Stream + Spring Cloud Sleuth integrates with http://projects.spring.io/spring-integration/[Spring Integration]. It creates spans for publish and subscribe events. To disable Spring Integration instrumentation, set `spring.sleuth.integration.enabled` to `false`. @@ -1137,11 +1139,19 @@ Decorating the Spring Integration Executor Channel with `TraceableExecutorServic ==== Spring RabbitMq -We instrument the `RabbiTemplate` so that tracing headers get injected +We instrument the `RabbitTemplate` so that tracing headers get injected into the message. To block this feature, set `spring.sleuth.messaging.enabled` to `false`. +==== Spring Kafka + +We instrument the Spring Kafka's `ProducerFactory` and `ConsumerFactory` +so that tracing headers get injected into the created Spring Kafka's +`Producer` and `Consumer`. + +To block this feature, set `spring.sleuth.messaging.enabled` to `false`. + === Zuul We instrument the Zuul Ribbon integration by enriching the Ribbon requests with tracing information. diff --git a/spring-cloud-sleuth-core/pom.xml b/spring-cloud-sleuth-core/pom.xml index c6c3e2124..33e36f35d 100644 --- a/spring-cloud-sleuth-core/pom.xml +++ b/spring-cloud-sleuth-core/pom.xml @@ -96,6 +96,11 @@ spring-rabbit true + + org.springframework.kafka + spring-kafka + true + org.springframework.security.oauth spring-security-oauth2 @@ -173,6 +178,10 @@ io.zipkin.brave brave-instrumentation-spring-rabbit + + io.zipkin.brave + brave-instrumentation-kafka-clients + io.zipkin.brave brave-instrumentation-httpclient diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfiguration.java index 1037ee083..fc6188cb4 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfiguration.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfiguration.java @@ -17,7 +17,14 @@ package org.springframework.cloud.sleuth.instrument.messaging; import brave.Tracing; +import brave.kafka.clients.KafkaTracing; import brave.spring.rabbit.SpringRabbitTracing; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.aspectj.lang.ProceedingJoinPoint; +import org.aspectj.lang.annotation.Around; +import org.aspectj.lang.annotation.Aspect; +import org.aspectj.lang.annotation.Pointcut; import org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.BeansException; @@ -32,6 +39,7 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties import org.springframework.cloud.sleuth.autoconfig.TraceAutoConfiguration; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.kafka.core.ProducerFactory; /** * {@link org.springframework.boot.autoconfigure.EnableAutoConfiguration @@ -67,6 +75,24 @@ public class TraceMessagingAutoConfiguration { return new SleuthRabbitBeanPostProcessor(beanFactory); } } + + @Configuration + @ConditionalOnClass(ProducerFactory.class) + protected static class SleuthKafkaConfiguration { + + @Bean + @ConditionalOnMissingBean + KafkaTracing kafkaTracing(Tracing tracing) { + return KafkaTracing.create(tracing); + } + + @Bean + // for tests + @ConditionalOnMissingBean + SleuthKafkaAspect sleuthKafkaAspect(KafkaTracing kafkaTracing) { + return new SleuthKafkaAspect(kafkaTracing); + } + } } class SleuthRabbitBeanPostProcessor implements BeanPostProcessor { @@ -96,4 +122,32 @@ class SleuthRabbitBeanPostProcessor implements BeanPostProcessor { } return this.tracing; } +} + +@Aspect +class SleuthKafkaAspect { + + private final KafkaTracing kafkaTracing; + + SleuthKafkaAspect(KafkaTracing kafkaTracing) { + this.kafkaTracing = kafkaTracing; + } + + @Pointcut("execution(public * org.springframework.kafka.core.ProducerFactory.createProducer(..))") + private void anyProducerFactory() { } // NOSONAR + + @Pointcut("execution(public * org.springframework.kafka.core.ConsumerFactory.createConsumer(..))") + private void anyConsumerFactory() { } // NOSONAR + + @Around("anyProducerFactory()") + public Object wrapProducerFactory(ProceedingJoinPoint pjp) throws Throwable { + Producer producer = (Producer) pjp.proceed(); + return this.kafkaTracing.producer(producer); + } + + @Around("anyConsumerFactory()") + public Object wrapConsumerFactory(ProceedingJoinPoint pjp) throws Throwable { + Consumer consumer = (Consumer) pjp.proceed(); + return this.kafkaTracing.consumer(consumer); + } } \ No newline at end of file diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java index bde8d9f59..3e654945d 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java @@ -16,9 +16,11 @@ package org.springframework.cloud.sleuth.instrument.messaging; +import brave.kafka.clients.KafkaTracing; import brave.sampler.Sampler; import brave.spring.rabbit.SpringRabbitTracing; import com.rabbitmq.client.Channel; +import org.aspectj.lang.ProceedingJoinPoint; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -41,6 +43,8 @@ import org.springframework.boot.test.mock.mockito.SpyBean; import org.springframework.cloud.sleuth.util.ArrayListSpanReporter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.kafka.core.ConsumerFactory; +import org.springframework.kafka.core.ProducerFactory; import org.springframework.test.context.junit4.SpringRunner; import static org.assertj.core.api.BDDAssertions.then; @@ -56,6 +60,9 @@ public class TraceMessagingAutoConfigurationTests { @Autowired RabbitTemplate rabbitTemplate; @Autowired ArrayListSpanReporter reporter; @Autowired TestSleuthRabbitBeanPostProcessor postProcessor; + @Autowired MySleuthKafkaAspect mySleuthKafkaAspect; + @Autowired ProducerFactory producerFactory; + @Autowired ConsumerFactory consumerFactory; @Test public void should_wrap_rabbit_template() { @@ -63,6 +70,15 @@ public class TraceMessagingAutoConfigurationTests { then(this.postProcessor.rabbitTracingCalled).isTrue(); } + @Test + public void should_wrap_kafka() { + this.producerFactory.createProducer(); + then(this.mySleuthKafkaAspect.producerWrapped).isTrue(); + + this.consumerFactory.createConsumer(); + then(this.mySleuthKafkaAspect.consumerWrapped).isTrue(); + } + @Configuration @EnableAutoConfiguration protected static class Config { @@ -77,6 +93,9 @@ public class TraceMessagingAutoConfigurationTests { @Bean SleuthRabbitBeanPostProcessor postProcessor(BeanFactory beanFactory) { return new TestSleuthRabbitBeanPostProcessor(beanFactory); } + @Bean SleuthKafkaAspect sleuthKafkaAspect(KafkaTracing kafkaTracing) { + return new MySleuthKafkaAspect(kafkaTracing); + } } } @@ -92,4 +111,26 @@ class TestSleuthRabbitBeanPostProcessor extends SleuthRabbitBeanPostProcessor { this.rabbitTracingCalled = true; return super.rabbitTracing(); } +} + +class MySleuthKafkaAspect extends SleuthKafkaAspect { + + boolean producerWrapped; + boolean consumerWrapped; + + MySleuthKafkaAspect(KafkaTracing kafkaTracing) { + super(kafkaTracing); + } + + @Override public Object wrapProducerFactory(ProceedingJoinPoint pjp) + throws Throwable { + this.producerWrapped = true; + return super.wrapProducerFactory(pjp); + } + + @Override public Object wrapConsumerFactory(ProceedingJoinPoint pjp) + throws Throwable { + this.consumerWrapped = true; + return super.wrapConsumerFactory(pjp); + } } \ No newline at end of file