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 862cd9e17..bab5700b3 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 @@ -61,6 +61,7 @@ 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.kafka.support.DefaultKafkaHeaderMapper; import org.springframework.util.ReflectionUtils; /** @@ -118,6 +119,21 @@ public class TraceMessagingAutoConfiguration { .build(); } + @Bean(name = SleuthDefaultKafkaHeaderMapper.BEAN_NAME) + @ConditionalOnMissingBean + DefaultKafkaHeaderMapper sleuthDefaultKafkaHeaderMapper() { + return new SleuthDefaultKafkaHeaderMapper(); + } + + @Bean + @ConditionalOnProperty(value = "spring.sleuth.messaging.kafka.mapper.enabled", + matchIfMissing = true) + // for tests + @ConditionalOnMissingBean + SleuthKafkaHeaderMapperBeanPostProcessor sleuthDefaultKafkaHeaderMapperBeanPostProcessor() { + return new SleuthKafkaHeaderMapperBeanPostProcessor(); + } + @Bean // for tests @ConditionalOnMissingBean @@ -375,3 +391,34 @@ class TracingJmsBeanPostProcessor implements BeanPostProcessor { } } + +class SleuthDefaultKafkaHeaderMapper extends DefaultKafkaHeaderMapper { + + // related to #1430 + static final String BEAN_NAME = "kafkaBinderHeaderMapper"; + + SleuthDefaultKafkaHeaderMapper() { + setMapAllStringsOut(true); + } + +} + +class SleuthKafkaHeaderMapperBeanPostProcessor implements BeanPostProcessor { + + @Override + public Object postProcessAfterInitialization(Object bean, String beanName) + throws BeansException { + if (bean instanceof SleuthDefaultKafkaHeaderMapper) { + return sleuthDefaultKafkaHeaderMapper(bean); + } + else if (bean instanceof DefaultKafkaHeaderMapper) { + ((DefaultKafkaHeaderMapper) bean).setMapAllStringsOut(true); + } + return bean; + } + + Object sleuthDefaultKafkaHeaderMapper(Object bean) { + return bean; + } + +} diff --git a/spring-cloud-sleuth-core/src/main/resources/META-INF/additional-spring-configuration-metadata.json b/spring-cloud-sleuth-core/src/main/resources/META-INF/additional-spring-configuration-metadata.json index 33b8369b6..041306568 100644 --- a/spring-cloud-sleuth-core/src/main/resources/META-INF/additional-spring-configuration-metadata.json +++ b/spring-cloud-sleuth-core/src/main/resources/META-INF/additional-spring-configuration-metadata.json @@ -53,6 +53,30 @@ "type": "java.lang.Boolean", "description": "Enable span information propagation when using Redis.", "defaultValue": true + }, + { + "name": "spring.sleuth.messaging.jms.enabled", + "type": "java.lang.Boolean", + "description": "Enable tracing of JMS.", + "defaultValue": true + }, + { + "name": "spring.sleuth.messaging.rabbit.enabled", + "type": "java.lang.Boolean", + "description": "Enable tracing of RabbitMQ.", + "defaultValue": true + }, + { + "name": "spring.sleuth.messaging.kafka.enabled", + "type": "java.lang.Boolean", + "description": "Enable tracing of Kafka.", + "defaultValue": true + }, + { + "name": "spring.sleuth.messaging.kafka.mapper.enabled", + "type": "java.lang.Boolean", + "description": "Enable DefaultKafkaHeaderMapper tracing for Kafka.", + "defaultValue": true } ] } diff --git a/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java b/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java index 731639cc9..b4b09a014 100644 --- a/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java +++ b/tests/spring-cloud-sleuth-instrumentation-messaging-tests/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TraceMessagingAutoConfigurationTests.java @@ -49,8 +49,7 @@ import static org.assertj.core.api.BDDAssertions.then; * @author Marcin Grzejszczak */ @RunWith(SpringRunner.class) -@SpringBootTest(classes = TraceMessagingAutoConfigurationTests.Config.class, - webEnvironment = SpringBootTest.WebEnvironment.NONE) +@SpringBootTest(classes = TraceMessagingAutoConfigurationTests.Config.class, webEnvironment = SpringBootTest.WebEnvironment.NONE) public class TraceMessagingAutoConfigurationTests { @Autowired @@ -68,6 +67,9 @@ public class TraceMessagingAutoConfigurationTests { @Autowired MySleuthKafkaAspect mySleuthKafkaAspect; + @Autowired + TestSleuthKafkaHeaderMapperBeanPostProcessor testSleuthKafkaHeaderMapperBeanPostProcessor; + @Autowired ProducerFactory producerFactory; @@ -95,6 +97,8 @@ public class TraceMessagingAutoConfigurationTests { then(this.mySleuthKafkaAspect.consumerWrapped).isTrue(); then(this.mySleuthKafkaAspect.adapterWrapped).isTrue(); + + then(this.testSleuthKafkaHeaderMapperBeanPostProcessor.tracingCalled).isTrue(); } @Configuration @@ -128,6 +132,11 @@ public class TraceMessagingAutoConfigurationTests { return new TestSleuthJmsBeanPostProcessor(beanFactory); } + @Bean + TestSleuthKafkaHeaderMapperBeanPostProcessor testSleuthKafkaHeaderMapperBeanPostProcessor() { + return new TestSleuthKafkaHeaderMapperBeanPostProcessor(); + } + @KafkaListener(topics = "backend", groupId = "foo") public void onMessage(ConsumerRecord message) { System.err.println(message); @@ -202,3 +211,15 @@ class TestSleuthJmsBeanPostProcessor extends TracingConnectionFactoryBeanPostPro } } + +class TestSleuthKafkaHeaderMapperBeanPostProcessor + extends SleuthKafkaHeaderMapperBeanPostProcessor { + + boolean tracingCalled = false; + + @Override + Object sleuthDefaultKafkaHeaderMapper(Object bean) { + this.tracingCalled = true; + return super.sleuthDefaultKafkaHeaderMapper(bean); + } +}