Removed the workaround since it got fixed in kafka; fixes gh-1515

This commit is contained in:
Marcin Grzejszczak
2020-01-09 09:24:55 +01:00
parent 73f762285e
commit 5a44385448
2 changed files with 0 additions and 70 deletions

View File

@@ -72,7 +72,6 @@ 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.lang.Nullable;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessagingException;
@@ -161,21 +160,6 @@ 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
@@ -457,37 +441,6 @@ 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;
}
}
class SqsQueueMessageHandlerFactory extends QueueMessageHandlerFactory {
private TracingMethodMessageHandlerAdapter handlerAdapter;

View File

@@ -75,9 +75,6 @@ public class TraceMessagingAutoConfigurationTests {
@Autowired
MySleuthKafkaAspect mySleuthKafkaAspect;
@Autowired
TestSleuthKafkaHeaderMapperBeanPostProcessor testSleuthKafkaHeaderMapperBeanPostProcessor;
@Autowired
ProducerFactory producerFactory;
@@ -105,8 +102,6 @@ public class TraceMessagingAutoConfigurationTests {
then(this.mySleuthKafkaAspect.consumerWrapped).isTrue();
then(this.mySleuthKafkaAspect.adapterWrapped).isTrue();
then(this.testSleuthKafkaHeaderMapperBeanPostProcessor.tracingCalled).isTrue();
}
@Test
@@ -189,11 +184,6 @@ 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);
@@ -269,19 +259,6 @@ class TestSleuthJmsBeanPostProcessor extends TracingConnectionFactoryBeanPostPro
}
class TestSleuthKafkaHeaderMapperBeanPostProcessor
extends SleuthKafkaHeaderMapperBeanPostProcessor {
boolean tracingCalled = false;
@Override
Object sleuthDefaultKafkaHeaderMapper(Object bean) {
this.tracingCalled = true;
return super.sleuthDefaultKafkaHeaderMapper(bean);
}
}
@Configuration
class ProducerSamplerConfig {