Added improved support of DefaultKafkaHeaderMapper wrapping for Kafka

without this change there was a propagation gap due to additional quotes in the headers
with this change we're wrapping the DefaultKafkaHeaderMapper to not add additional quotes

fixes gh-1430
This commit is contained in:
Marcin Grzejszczak
2019-10-02 13:13:37 +02:00
parent 3a2e76f588
commit 29c9afe403
3 changed files with 94 additions and 2 deletions

View File

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

View File

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

View File

@@ -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);
}
}