Added support for message containers; fixes gh-2093

This commit is contained in:
Marcin Grzejszczak
2022-01-03 13:57:01 +01:00
parent 03e498400b
commit abdf5754c2
9 changed files with 248 additions and 23 deletions

View File

@@ -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

View File

@@ -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<ConsumerRecord<?, ?>> extractor) {
return new TracingKafkaAspect(tracer, propagator, extractor);
}
}

View File

@@ -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<ConsumerPostProcessor> postProcessors = new ArrayList<>();

View File

@@ -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(..))")

View File

@@ -31,20 +31,25 @@ final class KafkaTracingUtils {
private KafkaTracingUtils() {
}
static <K, V> void buildAndFinishSpan(ConsumerRecord<K, V> consumerRecord, Propagator propagator,
Propagator.Getter<ConsumerRecord<?, ?>> 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 <K, V> void buildAndFinishSpan(SleuthKafkaSpan sleuthKafkaSpan, ConsumerRecord<K, V> consumerRecord,
Propagator propagator, Propagator.Getter<ConsumerRecord<?, ?>> extractor) {
Span span = buildSpan(sleuthKafkaSpan, consumerRecord, propagator, extractor);
if (log.isDebugEnabled()) {
log.debug("Extracted span from event headers " + span);
}
span.end();
}
static <K, V> Span buildSpan(SleuthKafkaSpan sleuthKafkaSpan, ConsumerRecord<K, V> consumerRecord,
Propagator propagator, Propagator.Getter<ConsumerRecord<?, ?>> 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();
}
}

View File

@@ -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.
*/

View File

@@ -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<ConsumerRecord<?, ?>> extractor;
public TracingKafkaAspect(Tracer tracer, Propagator propagator, Propagator.Getter<ConsumerRecord<?, ?>> 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();
}
}

View File

@@ -130,7 +130,8 @@ public class TracingKafkaConsumer<K, V> implements Consumer<K, V> {
public ConsumerRecords<K, V> poll(long l) {
ConsumerRecords<K, V> consumerRecords = this.delegate.poll(l);
for (ConsumerRecord<K, V> 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<K, V> implements Consumer<K, V> {
public ConsumerRecords<K, V> poll(Duration duration) {
ConsumerRecords<K, V> consumerRecords = this.delegate.poll(duration);
for (ConsumerRecord<K, V> consumerRecord : consumerRecords) {
KafkaTracingUtils.buildAndFinishSpan(consumerRecord, propagator(), extractor());
KafkaTracingUtils.buildAndFinishSpan(SleuthKafkaSpan.KAFKA_CONSUMER_SPAN, consumerRecord, propagator(),
extractor());
}
return consumerRecords;
}

View File

@@ -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<T extends MessageListener> implements MethodInterceptor {
private static final Log log = LogFactory.getLog(TracingMessageListenerMethodInterceptor.class);
private final Tracer tracer;
private final Propagator propagator;
private final Propagator.Getter<ConsumerRecord<?, ?>> extractor;
TracingMessageListenerMethodInterceptor(Tracer tracer, Propagator propagator,
Propagator.Getter<ConsumerRecord<?, ?>> 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;
}
}