Adds Messaging sampling (#1483)
Same as #1456 except for messaging. This completes remote span sampling!
This commit is contained in:
committed by
Marcin Grzejszczak
parent
848442e99c
commit
d1dddbf2c9
@@ -876,6 +876,34 @@ In the following example, we register the `TracingFilter` bean, add the `ZIPKIN-
|
||||
include::{project-root}/tests/spring-cloud-sleuth-instrumentation-mvc-tests/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceFilterIntegrationTests.java[tags=response_headers,indent=0]
|
||||
----
|
||||
|
||||
=== Messaging
|
||||
|
||||
Sleuth automatically configures the `MessagingTracing` bean which serves as a
|
||||
foundation for Messaging instrumentation such as Kafka or JMS.
|
||||
|
||||
If a customization of producer / consumer sampling of messaging traces is required,
|
||||
just register a bean of type `brave.sampler.SamplerFunction<MessagingRequest>` and
|
||||
name the bean `sleuthProducerSampler` for producer sampler and `sleuthConsumerSampler`
|
||||
for consumer sampler.
|
||||
|
||||
For your convenience the `@ProducerSampler` and `@ConsumerSampler`
|
||||
annotations can be used to inject the proper beans or to reference the bean
|
||||
names via their static String `NAME` fields.
|
||||
|
||||
Ex. Here's a sampler that traces 100 consumer requests per second, except for
|
||||
the "alerts" channel. Other requests will use a global rate provided by the
|
||||
`Tracing` component.
|
||||
|
||||
[source,java]
|
||||
----
|
||||
@Configuration
|
||||
class Config {
|
||||
include::{project-root}/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationIntegrationTests.java[tags=custom_messaging_server_sampler,indent=2]
|
||||
}
|
||||
----
|
||||
|
||||
For more, see https://github.com/openzipkin/brave/tree/master/instrumentation/messaging#sampling-policy
|
||||
|
||||
=== RPC
|
||||
|
||||
Sleuth automatically configures the `RpcTracing` bean which serves as a
|
||||
|
||||
@@ -211,6 +211,10 @@
|
||||
<groupId>io.zipkin.brave</groupId>
|
||||
<artifactId>brave-context-log4j2</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.zipkin.brave</groupId>
|
||||
<artifactId>brave-instrumentation-messaging</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.zipkin.brave</groupId>
|
||||
<artifactId>brave-instrumentation-rpc</artifactId>
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
/*
|
||||
* Copyright 2013-2019 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.messaging;
|
||||
|
||||
import java.lang.annotation.Documented;
|
||||
import java.lang.annotation.ElementType;
|
||||
import java.lang.annotation.Inherited;
|
||||
import java.lang.annotation.Retention;
|
||||
import java.lang.annotation.RetentionPolicy;
|
||||
import java.lang.annotation.Target;
|
||||
|
||||
import brave.sampler.SamplerFunction;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
|
||||
/**
|
||||
* Annotate a consumer {@link SamplerFunction} that should be injected to
|
||||
* {@link brave.messaging.MessagingTracing.Builder#consumerSampler(SamplerFunction)}.
|
||||
*
|
||||
* @since 2.2.0
|
||||
* @see Qualifier
|
||||
*/
|
||||
@Target({ ElementType.FIELD, ElementType.METHOD, ElementType.PARAMETER, ElementType.TYPE,
|
||||
ElementType.ANNOTATION_TYPE })
|
||||
@Retention(RetentionPolicy.RUNTIME)
|
||||
@Inherited
|
||||
@Documented
|
||||
@Qualifier(ConsumerSampler.NAME)
|
||||
public @interface ConsumerSampler {
|
||||
|
||||
/**
|
||||
* Default name for message consumer sampler.
|
||||
*/
|
||||
String NAME = "sleuthConsumerSampler";
|
||||
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
/*
|
||||
* Copyright 2013-2019 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.messaging;
|
||||
|
||||
import java.lang.annotation.Documented;
|
||||
import java.lang.annotation.ElementType;
|
||||
import java.lang.annotation.Inherited;
|
||||
import java.lang.annotation.Retention;
|
||||
import java.lang.annotation.RetentionPolicy;
|
||||
import java.lang.annotation.Target;
|
||||
|
||||
import brave.sampler.SamplerFunction;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Qualifier;
|
||||
|
||||
/**
|
||||
* Annotate a producer {@link SamplerFunction} that should be injected to
|
||||
* {@link brave.messaging.MessagingTracing.Builder#producerSampler(SamplerFunction)}.
|
||||
*
|
||||
* @since 2.2.0
|
||||
* @see Qualifier
|
||||
*/
|
||||
@Target({ ElementType.FIELD, ElementType.METHOD, ElementType.PARAMETER, ElementType.TYPE,
|
||||
ElementType.ANNOTATION_TYPE })
|
||||
@Retention(RetentionPolicy.RUNTIME)
|
||||
@Inherited
|
||||
@Documented
|
||||
@Qualifier(ProducerSampler.NAME)
|
||||
public @interface ProducerSampler {
|
||||
|
||||
/**
|
||||
* Default name for messaging producer sampler.
|
||||
*/
|
||||
String NAME = "sleuthProducerSampler";
|
||||
|
||||
}
|
||||
@@ -17,6 +17,7 @@
|
||||
package org.springframework.cloud.sleuth.instrument.messaging;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
|
||||
@@ -25,7 +26,11 @@ import brave.Tracer;
|
||||
import brave.Tracing;
|
||||
import brave.jms.JmsTracing;
|
||||
import brave.kafka.clients.KafkaTracing;
|
||||
import brave.propagation.Propagation;
|
||||
import brave.messaging.MessagingRequest;
|
||||
import brave.messaging.MessagingTracing;
|
||||
import brave.messaging.MessagingTracingCustomizer;
|
||||
import brave.propagation.Propagation.Getter;
|
||||
import brave.sampler.SamplerFunction;
|
||||
import brave.spring.rabbit.SpringRabbitTracing;
|
||||
import org.aopalliance.intercept.MethodInterceptor;
|
||||
import org.aopalliance.intercept.MethodInvocation;
|
||||
@@ -44,6 +49,7 @@ import org.springframework.amqp.rabbit.core.RabbitTemplate;
|
||||
import org.springframework.aop.framework.ProxyFactoryBean;
|
||||
import org.springframework.beans.BeansException;
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.beans.factory.config.BeanDefinition;
|
||||
import org.springframework.beans.factory.config.BeanPostProcessor;
|
||||
import org.springframework.boot.autoconfigure.AutoConfigureAfter;
|
||||
@@ -67,6 +73,7 @@ import org.springframework.kafka.listener.MessageListener;
|
||||
import org.springframework.kafka.listener.MessageListenerContainer;
|
||||
import org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter;
|
||||
import org.springframework.kafka.support.DefaultKafkaHeaderMapper;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessagingException;
|
||||
import org.springframework.messaging.converter.MessageConverter;
|
||||
@@ -89,6 +96,29 @@ import org.springframework.util.ReflectionUtils;
|
||||
@EnableConfigurationProperties(SleuthMessagingProperties.class)
|
||||
public class TraceMessagingAutoConfiguration {
|
||||
|
||||
@Autowired(required = false)
|
||||
List<MessagingTracingCustomizer> messagingTracingCustomizers = new ArrayList<>();
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
// NOTE: stable bean name as might be used outside sleuth
|
||||
MessagingTracing messagingTracing(Tracing tracing,
|
||||
@Nullable @ProducerSampler SamplerFunction<MessagingRequest> producerSampler,
|
||||
@Nullable @ConsumerSampler SamplerFunction<MessagingRequest> consumerSampler) {
|
||||
|
||||
MessagingTracing.Builder builder = MessagingTracing.newBuilder(tracing);
|
||||
if (producerSampler != null) {
|
||||
builder.producerSampler(producerSampler);
|
||||
}
|
||||
if (consumerSampler != null) {
|
||||
builder.consumerSampler(consumerSampler);
|
||||
}
|
||||
for (MessagingTracingCustomizer customizer : this.messagingTracingCustomizers) {
|
||||
customizer.customize(builder);
|
||||
}
|
||||
return builder.build();
|
||||
}
|
||||
|
||||
@Configuration(proxyBeanMethods = false)
|
||||
@ConditionalOnProperty(value = "spring.sleuth.messaging.rabbit.enabled",
|
||||
matchIfMissing = true)
|
||||
@@ -105,9 +135,9 @@ public class TraceMessagingAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
SpringRabbitTracing springRabbitTracing(Tracing tracing,
|
||||
SpringRabbitTracing springRabbitTracing(MessagingTracing messagingTracing,
|
||||
SleuthMessagingProperties properties) {
|
||||
return SpringRabbitTracing.newBuilder(tracing)
|
||||
return SpringRabbitTracing.newBuilder(messagingTracing)
|
||||
.remoteServiceName(
|
||||
properties.getMessaging().getRabbit().getRemoteServiceName())
|
||||
.build();
|
||||
@@ -123,8 +153,9 @@ public class TraceMessagingAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
KafkaTracing kafkaTracing(Tracing tracing, SleuthMessagingProperties properties) {
|
||||
return KafkaTracing.newBuilder(tracing)
|
||||
KafkaTracing kafkaTracing(MessagingTracing messagingTracing,
|
||||
SleuthMessagingProperties properties) {
|
||||
return KafkaTracing.newBuilder(messagingTracing)
|
||||
.remoteServiceName(
|
||||
properties.getMessaging().getKafka().getRemoteServiceName())
|
||||
.build();
|
||||
@@ -164,8 +195,9 @@ public class TraceMessagingAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
JmsTracing jmsTracing(Tracing tracing, SleuthMessagingProperties properties) {
|
||||
return JmsTracing.newBuilder(tracing)
|
||||
JmsTracing jmsTracing(MessagingTracing messagingTracing,
|
||||
SleuthMessagingProperties properties) {
|
||||
return JmsTracing.newBuilder(messagingTracing)
|
||||
.remoteServiceName(
|
||||
properties.getMessaging().getJms().getRemoteServiceName())
|
||||
.build();
|
||||
@@ -207,9 +239,9 @@ public class TraceMessagingAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
TracingMethodMessageHandlerAdapter tracingMethodMessageHandlerAdapter(
|
||||
Tracing tracing,
|
||||
Propagation.Getter<MessageHeaderAccessor, String> traceMessagePropagationGetter) {
|
||||
return new TracingMethodMessageHandlerAdapter(tracing,
|
||||
MessagingTracing messagingTracing,
|
||||
Getter<MessageHeaderAccessor, String> traceMessagePropagationGetter) {
|
||||
return new TracingMethodMessageHandlerAdapter(messagingTracing,
|
||||
traceMessagePropagationGetter);
|
||||
}
|
||||
|
||||
@@ -477,7 +509,7 @@ class SqsQueueMessageHandlerFactory extends QueueMessageHandlerFactory {
|
||||
class SqsQueueMessageHandler extends QueueMessageHandler {
|
||||
|
||||
// copied from QueueMessageHandler
|
||||
private static final String LOGICAL_RESOURCE_ID = "LogicalResourceId";
|
||||
static final String LOGICAL_RESOURCE_ID = "LogicalResourceId";
|
||||
|
||||
private TracingMethodMessageHandlerAdapter handlerAdapter;
|
||||
|
||||
|
||||
@@ -21,8 +21,10 @@ import java.util.function.BiConsumer;
|
||||
import brave.Span;
|
||||
import brave.Tracer;
|
||||
import brave.Tracing;
|
||||
import brave.propagation.Propagation;
|
||||
import brave.propagation.TraceContext;
|
||||
import brave.messaging.ConsumerRequest;
|
||||
import brave.messaging.MessagingTracing;
|
||||
import brave.propagation.Propagation.Getter;
|
||||
import brave.propagation.TraceContext.Extractor;
|
||||
import brave.propagation.TraceContextOrSamplingFlags;
|
||||
|
||||
import org.springframework.messaging.Message;
|
||||
@@ -30,6 +32,7 @@ import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.support.MessageHeaderAccessor;
|
||||
|
||||
import static brave.Span.Kind.CONSUMER;
|
||||
import static org.springframework.cloud.sleuth.instrument.messaging.SqsQueueMessageHandler.LOGICAL_RESOURCE_ID;
|
||||
|
||||
/**
|
||||
* Adds tracing extraction to an instance of
|
||||
@@ -46,22 +49,26 @@ import static brave.Span.Kind.CONSUMER;
|
||||
*/
|
||||
class TracingMethodMessageHandlerAdapter {
|
||||
|
||||
private Tracing tracing;
|
||||
private final Tracing tracing;
|
||||
|
||||
private Tracer tracer;
|
||||
private final Tracer tracer;
|
||||
|
||||
private TraceContext.Extractor<MessageHeaderAccessor> extractor;
|
||||
private final Extractor<MessageConsumerRequest> extractor;
|
||||
|
||||
TracingMethodMessageHandlerAdapter(Tracing tracing,
|
||||
Propagation.Getter<MessageHeaderAccessor, String> traceMessagePropagationGetter) {
|
||||
this.tracing = tracing;
|
||||
private final Getter<MessageHeaderAccessor, String> getter;
|
||||
|
||||
TracingMethodMessageHandlerAdapter(MessagingTracing messagingTracing,
|
||||
Getter<MessageHeaderAccessor, String> getter) {
|
||||
this.tracing = messagingTracing.tracing();
|
||||
this.tracer = tracing.tracer();
|
||||
this.extractor = tracing.propagation().extractor(traceMessagePropagationGetter);
|
||||
this.extractor = tracing.propagation().extractor(MessageConsumerRequest.GETTER);
|
||||
this.getter = getter;
|
||||
}
|
||||
|
||||
void wrapMethodMessageHandler(Message<?> message, MessageHandler messageHandler,
|
||||
BiConsumer<Span, Message<?>> messageSpanTagger) {
|
||||
TraceContextOrSamplingFlags extracted = extractAndClearHeaders(message);
|
||||
MessageConsumerRequest request = new MessageConsumerRequest(message, this.getter);
|
||||
TraceContextOrSamplingFlags extracted = extractAndClearHeaders(request);
|
||||
|
||||
Span consumerSpan = tracer.nextSpan(extracted);
|
||||
Span listenerSpan = tracer.newChild(consumerSpan.context());
|
||||
@@ -95,15 +102,77 @@ class TracingMethodMessageHandlerAdapter {
|
||||
}
|
||||
}
|
||||
|
||||
private TraceContextOrSamplingFlags extractAndClearHeaders(Message<?> message) {
|
||||
MessageHeaderAccessor headers = MessageHeaderAccessor.getMutableAccessor(message);
|
||||
TraceContextOrSamplingFlags extracted = extractor.extract(headers);
|
||||
private TraceContextOrSamplingFlags extractAndClearHeaders(
|
||||
MessageConsumerRequest request) {
|
||||
TraceContextOrSamplingFlags extracted = extractor.extract(request);
|
||||
|
||||
for (String propagationKey : tracing.propagation().keys()) {
|
||||
headers.removeHeader(propagationKey);
|
||||
request.removeHeader(propagationKey);
|
||||
}
|
||||
|
||||
return extracted;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
final class MessageConsumerRequest extends ConsumerRequest {
|
||||
|
||||
static final Getter<MessageConsumerRequest, String> GETTER = new Getter<MessageConsumerRequest, String>() {
|
||||
@Override
|
||||
public String get(MessageConsumerRequest request, String name) {
|
||||
return request.getHeader(name);
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "MessageConsumerRequest::getHeader";
|
||||
}
|
||||
};
|
||||
|
||||
final Message delegate;
|
||||
|
||||
final MessageHeaderAccessor mutableHeaders;
|
||||
|
||||
final Getter<MessageHeaderAccessor, String> getter;
|
||||
|
||||
MessageConsumerRequest(Message delegate,
|
||||
Getter<MessageHeaderAccessor, String> getter) {
|
||||
this.delegate = delegate;
|
||||
this.mutableHeaders = MessageHeaderAccessor.getMutableAccessor(delegate);
|
||||
this.getter = getter;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Span.Kind spanKind() {
|
||||
return Span.Kind.CONSUMER;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object unwrap() {
|
||||
return this.delegate;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String operation() {
|
||||
return "receive";
|
||||
}
|
||||
|
||||
@Override
|
||||
public String channelKind() {
|
||||
return "queue";
|
||||
}
|
||||
|
||||
@Override
|
||||
public String channelName() {
|
||||
return this.delegate.getHeaders().get(LOGICAL_RESOURCE_ID).toString();
|
||||
}
|
||||
|
||||
String getHeader(String name) {
|
||||
return this.getter.get(this.mutableHeaders, name);
|
||||
}
|
||||
|
||||
void removeHeader(String name) {
|
||||
this.mutableHeaders.removeHeader(name);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -18,8 +18,10 @@ package org.springframework.cloud.sleuth.autoconfig;
|
||||
|
||||
import brave.TracingCustomizer;
|
||||
import brave.http.HttpTracingCustomizer;
|
||||
import brave.messaging.MessagingTracingCustomizer;
|
||||
import brave.propagation.CurrentTraceContextCustomizer;
|
||||
import brave.propagation.ExtraFieldCustomizer;
|
||||
import brave.propagation.Propagation;
|
||||
import brave.rpc.RpcTracingCustomizer;
|
||||
import brave.sampler.Sampler;
|
||||
import org.junit.Test;
|
||||
@@ -27,11 +29,13 @@ import org.junit.Test;
|
||||
import org.springframework.boot.autoconfigure.AutoConfigurations;
|
||||
import org.springframework.boot.test.context.assertj.AssertableApplicationContext;
|
||||
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
|
||||
import org.springframework.cloud.sleuth.instrument.messaging.TraceMessagingAutoConfiguration;
|
||||
import org.springframework.cloud.sleuth.instrument.rpc.TraceRpcAutoConfiguration;
|
||||
import org.springframework.cloud.sleuth.instrument.web.TraceHttpAutoConfiguration;
|
||||
import org.springframework.cloud.sleuth.instrument.web.TraceWebAutoConfiguration;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.messaging.support.MessageHeaderAccessor;
|
||||
|
||||
import static org.assertj.core.api.BDDAssertions.then;
|
||||
|
||||
@@ -40,7 +44,9 @@ public class TraceAutoConfigurationCustomizersTests {
|
||||
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
|
||||
.withConfiguration(AutoConfigurations.of(TraceAutoConfiguration.class,
|
||||
TraceWebAutoConfiguration.class, TraceHttpAutoConfiguration.class,
|
||||
TraceRpcAutoConfiguration.class))
|
||||
TraceRpcAutoConfiguration.class,
|
||||
FakeSpringMessagingAutoConfiguration.class,
|
||||
TraceMessagingAutoConfiguration.class))
|
||||
.withUserConfiguration(Customizers.class);
|
||||
|
||||
@Test
|
||||
@@ -75,6 +81,17 @@ public class TraceAutoConfigurationCustomizersTests {
|
||||
then(bean.rpcCustomizerApplied).isTrue();
|
||||
}
|
||||
|
||||
// SQS has a dependency on the getter and this is better than exposing things public
|
||||
@Configuration
|
||||
static class FakeSpringMessagingAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
Propagation.Getter<MessageHeaderAccessor, String> traceMessagePropagationGetter() {
|
||||
return (headers, key) -> null;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Configuration
|
||||
static class Customizers {
|
||||
|
||||
@@ -88,6 +105,8 @@ public class TraceAutoConfigurationCustomizersTests {
|
||||
|
||||
boolean rpcCustomizerApplied;
|
||||
|
||||
boolean messagingCustomizerApplied;
|
||||
|
||||
@Bean
|
||||
TracingCustomizer sleuthTracingCustomizer() {
|
||||
return builder -> tracingCustomizerApplied = true;
|
||||
@@ -108,6 +127,11 @@ public class TraceAutoConfigurationCustomizersTests {
|
||||
return builder -> httpCustomizerApplied = true;
|
||||
}
|
||||
|
||||
@Bean
|
||||
MessagingTracingCustomizer sleuthMessagingCustomizer() {
|
||||
return builder -> messagingCustomizerApplied = true;
|
||||
}
|
||||
|
||||
@Bean
|
||||
RpcTracingCustomizer sleuthRpcTracingCustomizer() {
|
||||
return builder -> rpcCustomizerApplied = true;
|
||||
|
||||
@@ -0,0 +1,74 @@
|
||||
/*
|
||||
* Copyright 2013-2019 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.messaging;
|
||||
|
||||
import brave.messaging.MessagingRequest;
|
||||
import brave.messaging.MessagingRuleSampler;
|
||||
import brave.sampler.Matchers;
|
||||
import brave.sampler.RateLimitingSampler;
|
||||
import brave.sampler.Sampler;
|
||||
import brave.sampler.SamplerFunction;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.cloud.sleuth.util.ArrayListSpanReporter;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.test.context.junit4.SpringRunner;
|
||||
|
||||
import static brave.messaging.MessagingRequestMatchers.channelNameEquals;
|
||||
import static org.assertj.core.api.BDDAssertions.then;
|
||||
|
||||
@RunWith(SpringRunner.class)
|
||||
@SpringBootTest(classes = TraceMessagingAutoConfigurationIntegrationTests.Config.class,
|
||||
webEnvironment = SpringBootTest.WebEnvironment.NONE)
|
||||
public class TraceMessagingAutoConfigurationIntegrationTests {
|
||||
|
||||
@Autowired
|
||||
@ConsumerSampler
|
||||
SamplerFunction<MessagingRequest> sampler;
|
||||
|
||||
@Test
|
||||
public void should_inject_messaging_sampler() {
|
||||
then(this.sampler).isNotNull();
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@Configuration
|
||||
public static class Config {
|
||||
|
||||
@Bean
|
||||
ArrayListSpanReporter reporter() {
|
||||
return new ArrayListSpanReporter();
|
||||
}
|
||||
|
||||
// tag::custom_messaging_consumer_sampler[]
|
||||
@Bean(name = ConsumerSampler.NAME)
|
||||
SamplerFunction<MessagingRequest> myMessagingSampler() {
|
||||
return MessagingRuleSampler.newBuilder()
|
||||
.putRule(channelNameEquals("alerts"), Sampler.NEVER_SAMPLE)
|
||||
.putRule(Matchers.alwaysMatch(), RateLimitingSampler.create(100))
|
||||
.build();
|
||||
}
|
||||
// end::custom_messaging_consumer_sampler[]
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
@@ -31,8 +31,8 @@
|
||||
<name>spring-cloud-sleuth-dependencies</name>
|
||||
<description>Spring Cloud Sleuth Dependencies</description>
|
||||
<properties>
|
||||
<brave.version>5.8.0</brave.version>
|
||||
<brave.opentracing.version>0.34.1</brave.opentracing.version>
|
||||
<brave.version>5.9.0</brave.version>
|
||||
<brave.opentracing.version>0.35.0</brave.opentracing.version>
|
||||
<grpc.spring.boot.version>3.4.1</grpc.spring.boot.version>
|
||||
</properties>
|
||||
<dependencyManagement>
|
||||
|
||||
@@ -18,7 +18,11 @@ package org.springframework.cloud.sleuth.instrument.messaging;
|
||||
|
||||
import brave.Tracer;
|
||||
import brave.kafka.clients.KafkaTracing;
|
||||
import brave.messaging.MessagingRequest;
|
||||
import brave.messaging.MessagingTracing;
|
||||
import brave.sampler.Sampler;
|
||||
import brave.sampler.SamplerFunction;
|
||||
import brave.sampler.SamplerFunctions;
|
||||
import brave.spring.rabbit.SpringRabbitTracing;
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
@@ -32,8 +36,11 @@ import org.springframework.amqp.rabbit.core.RabbitTemplate;
|
||||
import org.springframework.beans.BeansException;
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.AutoConfigurations;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
|
||||
import org.springframework.cloud.sleuth.autoconfig.TraceAutoConfiguration;
|
||||
import org.springframework.cloud.sleuth.util.ArrayListSpanReporter;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
@@ -102,6 +109,55 @@ public class TraceMessagingAutoConfigurationTests {
|
||||
then(this.testSleuthKafkaHeaderMapperBeanPostProcessor.tracingCalled).isTrue();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void defaultsToBraveProducerSampler() {
|
||||
contextRunner().run((context) -> {
|
||||
SamplerFunction<MessagingRequest> producerSampler = context
|
||||
.getBean(MessagingTracing.class).producerSampler();
|
||||
|
||||
then(producerSampler).isSameAs(SamplerFunctions.deferDecision());
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
public void configuresUserProvidedProducerSampler() {
|
||||
contextRunner().withUserConfiguration(ProducerSamplerConfig.class)
|
||||
.run((context) -> {
|
||||
SamplerFunction<MessagingRequest> producerSampler = context
|
||||
.getBean(MessagingTracing.class).producerSampler();
|
||||
|
||||
then(producerSampler).isSameAs(ProducerSamplerConfig.INSTANCE);
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
public void defaultsToBraveConsumerSampler() {
|
||||
contextRunner().run((context) -> {
|
||||
SamplerFunction<MessagingRequest> consumerSampler = context
|
||||
.getBean(MessagingTracing.class).consumerSampler();
|
||||
|
||||
then(consumerSampler).isSameAs(SamplerFunctions.deferDecision());
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
public void configuresUserProvidedConsumerSampler() {
|
||||
contextRunner().withUserConfiguration(ConsumerSamplerConfig.class)
|
||||
.run((context) -> {
|
||||
SamplerFunction<MessagingRequest> consumerSampler = context
|
||||
.getBean(MessagingTracing.class).consumerSampler();
|
||||
|
||||
then(consumerSampler).isSameAs(ConsumerSamplerConfig.INSTANCE);
|
||||
});
|
||||
}
|
||||
|
||||
private ApplicationContextRunner contextRunner(String... propertyValues) {
|
||||
return new ApplicationContextRunner().withPropertyValues(propertyValues)
|
||||
.withConfiguration(AutoConfigurations.of(TraceAutoConfiguration.class,
|
||||
TraceMessagingAutoConfiguration.class,
|
||||
TraceMessagingAutoConfiguration.class));
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableAutoConfiguration
|
||||
protected static class Config {
|
||||
@@ -225,3 +281,27 @@ class TestSleuthKafkaHeaderMapperBeanPostProcessor
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Configuration
|
||||
class ProducerSamplerConfig {
|
||||
|
||||
static final SamplerFunction<MessagingRequest> INSTANCE = request -> null;
|
||||
|
||||
@Bean(ProducerSampler.NAME)
|
||||
SamplerFunction<MessagingRequest> sleuthProducerSampler() {
|
||||
return INSTANCE;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Configuration
|
||||
class ConsumerSamplerConfig {
|
||||
|
||||
static final SamplerFunction<MessagingRequest> INSTANCE = request -> null;
|
||||
|
||||
@Bean(ConsumerSampler.NAME)
|
||||
SamplerFunction<MessagingRequest> sleuthConsumerSampler() {
|
||||
return INSTANCE;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user