diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/brave/instrument/messaging/BraveMessagingAutoConfiguration.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/brave/instrument/messaging/BraveMessagingAutoConfiguration.java index 191d1350c..69e343d95 100644 --- a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/brave/instrument/messaging/BraveMessagingAutoConfiguration.java +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/brave/instrument/messaging/BraveMessagingAutoConfiguration.java @@ -114,7 +114,7 @@ public class BraveMessagingAutoConfiguration { @Configuration(proxyBeanMethods = false) @ConditionalOnProperty(value = "spring.sleuth.messaging.kafka.enabled", matchIfMissing = true) - @ConditionalOnClass(ProducerFactory.class) + @ConditionalOnClass({ KafkaTracing.class, ProducerFactory.class }) protected static class SleuthKafkaConfiguration { @Bean diff --git a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/SpringKafkaAutoConfiguration.java b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/SpringKafkaAutoConfiguration.java index 415e4a926..dc02dd8ec 100644 --- a/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/SpringKafkaAutoConfiguration.java +++ b/spring-cloud-sleuth-autoconfigure/src/main/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/SpringKafkaAutoConfiguration.java @@ -16,6 +16,8 @@ package org.springframework.cloud.sleuth.autoconfig.instrument.kafka; +import org.apache.kafka.clients.consumer.ConsumerRecord; + import org.springframework.beans.factory.BeanFactory; import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; @@ -24,6 +26,8 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingClas 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.TracingKafkaAspect; +import org.springframework.cloud.sleuth.propagation.Propagator; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.core.ProducerFactory; @@ -49,4 +53,10 @@ public class SpringKafkaAutoConfiguration { return new SpringKafkaFactoryBeanPostProcessor(beanFactory); } + @Bean + TracingKafkaAspect tracingKafkaAspect(Tracer tracer, Propagator propagator, + Propagator.Getter> extractor) { + return new TracingKafkaAspect(tracer, propagator, extractor); + } + } diff --git a/spring-cloud-sleuth-autoconfigure/src/test/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/SpringKafkaAutoConfigurationTests.java b/spring-cloud-sleuth-autoconfigure/src/test/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/SpringKafkaAutoConfigurationTests.java index 62f3d8582..c4ded880d 100644 --- a/spring-cloud-sleuth-autoconfigure/src/test/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/SpringKafkaAutoConfigurationTests.java +++ b/spring-cloud-sleuth-autoconfigure/src/test/java/org/springframework/cloud/sleuth/autoconfig/instrument/kafka/SpringKafkaAutoConfigurationTests.java @@ -28,6 +28,7 @@ import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.test.context.FilteredClassLoader; import org.springframework.boot.test.context.runner.ApplicationContextRunner; import org.springframework.cloud.sleuth.autoconfig.TraceNoOpAutoConfiguration; +import org.springframework.cloud.sleuth.instrument.kafka.TracingKafkaAspect; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.ConsumerPostProcessor; import org.springframework.kafka.core.ProducerFactory; @@ -44,8 +45,8 @@ class SpringKafkaAutoConfigurationTests { @Test void should_be_disabled_when_brave_on_classpath() { - this.contextRunner - .run((context) -> assertThat(context).doesNotHaveBean(SpringKafkaFactoryBeanPostProcessor.class)); + this.contextRunner.run((context) -> assertThat(context) + .doesNotHaveBean(SpringKafkaFactoryBeanPostProcessor.class).doesNotHaveBean(TracingKafkaAspect.class)); } @Test @@ -66,6 +67,12 @@ class SpringKafkaAutoConfigurationTests { .stream().filter(p -> p instanceof SpringKafkaConsumerPostProcessor).count() == 1)); } + @Test + void should_register_tracing_kafka_aspect() { + this.contextRunner.withClassLoader(new FilteredClassLoader(KafkaTracing.class)) + .run((context) -> assertThat(context).hasSingleBean(TracingKafkaAspect.class)); + } + class TestConsumerFactory implements ConsumerFactory { List postProcessors = new ArrayList<>(); diff --git a/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/instrument/messaging/SleuthKafkaAspect.java b/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/instrument/messaging/SleuthKafkaAspect.java index ad46acea5..86a0cac0d 100644 --- a/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/instrument/messaging/SleuthKafkaAspect.java +++ b/spring-cloud-sleuth-brave/src/main/java/org/springframework/cloud/sleuth/brave/instrument/messaging/SleuthKafkaAspect.java @@ -16,8 +16,6 @@ package org.springframework.cloud.sleuth.brave.instrument.messaging; -import java.lang.reflect.Field; - import brave.Tracer; import brave.kafka.clients.KafkaTracing; import org.apache.commons.logging.Log; @@ -31,8 +29,6 @@ import org.springframework.aop.framework.ProxyFactoryBean; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.MessageListener; import org.springframework.kafka.listener.MessageListenerContainer; -import org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter; -import org.springframework.util.ReflectionUtils; /** * Instruments Kafka related components. @@ -45,8 +41,6 @@ public class SleuthKafkaAspect { private static final Log log = LogFactory.getLog(SleuthKafkaAspect.class); - final Field recordMessageConverter; - private final KafkaTracing kafkaTracing; private final Tracer tracer; @@ -54,8 +48,6 @@ public class SleuthKafkaAspect { public SleuthKafkaAspect(KafkaTracing kafkaTracing, Tracer tracer) { this.kafkaTracing = kafkaTracing; this.tracer = tracer; - this.recordMessageConverter = ReflectionUtils.findField(MessagingMessageListenerAdapter.class, - "recordMessageConverter"); } @Pointcut("execution(public * org.springframework.kafka.config.KafkaListenerContainerFactory.createListenerContainer(..))") diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTracingUtils.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTracingUtils.java index 20e7a6858..54d632029 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTracingUtils.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/KafkaTracingUtils.java @@ -31,20 +31,25 @@ final class KafkaTracingUtils { private KafkaTracingUtils() { } - static void buildAndFinishSpan(ConsumerRecord consumerRecord, Propagator propagator, - Propagator.Getter> extractor) { - // @formatter:off - Span.Builder spanBuilder = AssertingSpanBuilder.of(SleuthKafkaSpan.KAFKA_CONSUMER_SPAN, propagator.extract(consumerRecord, extractor).kind(Span.Kind.CONSUMER)) - .name(SleuthKafkaSpan.KAFKA_CONSUMER_SPAN.getName()) - .tag(SleuthKafkaSpan.ConsumerTags.TOPIC, consumerRecord.topic()) - .tag(SleuthKafkaSpan.ConsumerTags.OFFSET, Long.toString(consumerRecord.offset())) - .tag(SleuthKafkaSpan.ConsumerTags.PARTITION, Integer.toString(consumerRecord.partition())); - // @formatter:on - Span span = spanBuilder.start(); + static void buildAndFinishSpan(SleuthKafkaSpan sleuthKafkaSpan, ConsumerRecord consumerRecord, + Propagator propagator, Propagator.Getter> extractor) { + Span span = buildSpan(sleuthKafkaSpan, consumerRecord, propagator, extractor); if (log.isDebugEnabled()) { log.debug("Extracted span from event headers " + span); } span.end(); } + static Span buildSpan(SleuthKafkaSpan sleuthKafkaSpan, ConsumerRecord consumerRecord, + Propagator propagator, Propagator.Getter> extractor) { + // @formatter:off + Span.Builder spanBuilder = AssertingSpanBuilder.of(sleuthKafkaSpan, propagator.extract(consumerRecord, extractor).kind(Span.Kind.CONSUMER)) + .name(sleuthKafkaSpan.getName()) + .tag(SleuthKafkaSpan.ConsumerTags.TOPIC, consumerRecord.topic()) + .tag(SleuthKafkaSpan.ConsumerTags.OFFSET, Long.toString(consumerRecord.offset())) + .tag(SleuthKafkaSpan.ConsumerTags.PARTITION, Integer.toString(consumerRecord.partition())); + // @formatter:on + return spanBuilder.start(); + } + } diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/SleuthKafkaSpan.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/SleuthKafkaSpan.java index e19cf3c54..5e2765a48 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/SleuthKafkaSpan.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/SleuthKafkaSpan.java @@ -41,6 +41,26 @@ enum SleuthKafkaSpan implements DocumentedSpan { } }, + /** + * Span created on the Kafka consumer side when using a MessageListener. + */ + KAFKA_ON_MESSAGE_SPAN { + @Override + public String getName() { + return "kafka.on-message"; + } + + @Override + public TagKey[] getTagKeys() { + return ConsumerTags.values(); + } + + @Override + public String prefix() { + return "kafka."; + } + }, + /** * Span created on the Kafka consumer side. */ diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaAspect.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaAspect.java new file mode 100644 index 000000000..b7d79cd85 --- /dev/null +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingKafkaAspect.java @@ -0,0 +1,102 @@ +/* + * 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.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.apache.kafka.clients.consumer.ConsumerRecord; +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.aop.framework.ProxyFactoryBean; +import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.kafka.listener.AbstractMessageListenerContainer; +import org.springframework.kafka.listener.MessageListener; +import org.springframework.kafka.listener.MessageListenerContainer; + +/** + * Instruments Kafka related components. + * + * @since 3.1.1 + * @author Marcin Grzejszczak + */ +@Aspect +public class TracingKafkaAspect { + + private static final Log log = LogFactory.getLog(TracingKafkaAspect.class); + + private final Tracer tracer; + + private final Propagator propagator; + + private final Propagator.Getter> extractor; + + public TracingKafkaAspect(Tracer tracer, Propagator propagator, Propagator.Getter> extractor) { + this.tracer = tracer; + this.propagator = propagator; + this.extractor = extractor; + } + + @Pointcut("execution(public * org.springframework.kafka.config.KafkaListenerContainerFactory.createListenerContainer(..))") + private void anyCreateListenerContainer() { + } // NOSONAR + + @Pointcut("execution(public * org.springframework.kafka.config.KafkaListenerContainerFactory.createContainer(..))") + private void anyCreateContainer() { + } // NOSONAR + + @Around("anyCreateListenerContainer() || anyCreateContainer()") + public Object wrapListenerContainerCreation(ProceedingJoinPoint pjp) throws Throwable { + MessageListenerContainer listener = (MessageListenerContainer) pjp.proceed(); + if (listener instanceof AbstractMessageListenerContainer) { + AbstractMessageListenerContainer container = (AbstractMessageListenerContainer) listener; + Object someMessageListener = container.getContainerProperties().getMessageListener(); + if (someMessageListener == null) { + if (log.isDebugEnabled()) { + log.debug("No message listener to wrap. Proceeding"); + } + } + else if (someMessageListener instanceof MessageListener) { + container.setupMessageListener(createProxy(someMessageListener)); + } + else { + if (log.isDebugEnabled()) { + log.debug("ATM we don't support Batch message listeners"); + } + } + } + else { + if (log.isDebugEnabled()) { + log.debug("Can't wrap this listener. Proceeding"); + } + } + return listener; + } + + @SuppressWarnings("unchecked") + Object createProxy(Object bean) { + ProxyFactoryBean factory = new ProxyFactoryBean(); + factory.setProxyTargetClass(true); + factory.addAdvice(new TracingMessageListenerMethodInterceptor(this.tracer, propagator, extractor)); + factory.setTarget(bean); + return factory.getObject(); + } + +} 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 d526a3cb5..719f49a51 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 @@ -130,7 +130,8 @@ public class TracingKafkaConsumer implements Consumer { public ConsumerRecords poll(long l) { ConsumerRecords consumerRecords = this.delegate.poll(l); for (ConsumerRecord consumerRecord : consumerRecords) { - KafkaTracingUtils.buildAndFinishSpan(consumerRecord, propagator(), extractor()); + KafkaTracingUtils.buildAndFinishSpan(SleuthKafkaSpan.KAFKA_CONSUMER_SPAN, consumerRecord, propagator(), + extractor()); } return consumerRecords; } @@ -139,7 +140,8 @@ public class TracingKafkaConsumer implements Consumer { public ConsumerRecords poll(Duration duration) { ConsumerRecords consumerRecords = this.delegate.poll(duration); for (ConsumerRecord consumerRecord : consumerRecords) { - KafkaTracingUtils.buildAndFinishSpan(consumerRecord, propagator(), extractor()); + KafkaTracingUtils.buildAndFinishSpan(SleuthKafkaSpan.KAFKA_CONSUMER_SPAN, consumerRecord, propagator(), + extractor()); } return consumerRecords; } diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingMessageListenerMethodInterceptor.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingMessageListenerMethodInterceptor.java new file mode 100644 index 000000000..09b920166 --- /dev/null +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/kafka/TracingMessageListenerMethodInterceptor.java @@ -0,0 +1,87 @@ +/* + * 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.aopalliance.intercept.MethodInterceptor; +import org.aopalliance.intercept.MethodInvocation; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.apache.kafka.clients.consumer.ConsumerRecord; + +import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.propagation.Propagator; +import org.springframework.kafka.listener.MessageListener; + +class TracingMessageListenerMethodInterceptor implements MethodInterceptor { + + private static final Log log = LogFactory.getLog(TracingMessageListenerMethodInterceptor.class); + + private final Tracer tracer; + + private final Propagator propagator; + + private final Propagator.Getter> extractor; + + TracingMessageListenerMethodInterceptor(Tracer tracer, Propagator propagator, + Propagator.Getter> extractor) { + this.tracer = tracer; + this.propagator = propagator; + this.extractor = extractor; + } + + @Override + public Object invoke(MethodInvocation invocation) throws Throwable { + if (!"onMessage".equals(invocation.getMethod().getName())) { + return invocation.proceed(); + } + Object[] arguments = invocation.getArguments(); + Object record = record(arguments); + if (record == null) { + return invocation.proceed(); + } + if (log.isDebugEnabled()) { + log.debug("Wrapping onMessage call"); + } + Span span = KafkaTracingUtils.buildSpan(SleuthKafkaSpan.KAFKA_ON_MESSAGE_SPAN, (ConsumerRecord) record, + this.propagator, this.extractor); + try (Tracer.SpanInScope ws = this.tracer.withSpan(span)) { + return invocation.proceed(); + } + catch (RuntimeException | Error e) { + String message = e.getMessage(); + if (message == null) { + message = e.getClass().getSimpleName(); + } + span.tag("error", message); + throw e; + } + finally { + span.end(); + } + } + + private Object record(Object[] arguments) { + for (Object object : arguments) { + if (object instanceof ConsumerRecord) { + return object; + } + } + return null; + } + +}