From 51c84c4b82444541d11be9d12d03c09151a10ee2 Mon Sep 17 00:00:00 2001 From: Dimitri Penner Date: Fri, 5 Jan 2018 15:44:53 +0100 Subject: [PATCH] GH-527: @KafkaListener: populate converter beans Fixes https://github.com/spring-projects/spring-kafka/issues/527 * Add test * Changed listener factory creation order * Add client id to manualConsumerFactory call * Refactor to add Converter beans to DefaultFormattingConversionService * Move addFormatters to afterSingletonsInstantiated * Add documentation * Improve documentation --- ...kaListenerAnnotationBeanPostProcessor.java | 41 ++++++++- .../EnableKafkaIntegrationTests.java | 92 ++++++++++++++++++- src/reference/asciidoc/kafka.adoc | 14 +++ 3 files changed, 140 insertions(+), 7 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java index e1201092..2600e14a 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2017 the original author or authors. + * Copyright 2014-2018 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. @@ -50,6 +50,10 @@ import org.springframework.context.expression.StandardBeanExpressionResolver; import org.springframework.core.MethodIntrospector; import org.springframework.core.Ordered; import org.springframework.core.annotation.AnnotationUtils; +import org.springframework.core.convert.converter.Converter; +import org.springframework.core.convert.converter.GenericConverter; +import org.springframework.format.Formatter; +import org.springframework.format.FormatterRegistry; import org.springframework.format.support.DefaultFormattingConversionService; import org.springframework.kafka.config.KafkaListenerConfigUtils; import org.springframework.kafka.config.KafkaListenerContainerFactory; @@ -98,6 +102,7 @@ import org.springframework.util.StringUtils; * @author Artem Bilan * @author Dariusz Szablinski * @author Venil Noronha + * @author Dimitri Penner * * @see KafkaListener * @see KafkaListenerErrorHandler @@ -224,6 +229,9 @@ public class KafkaListenerAnnotationBeanPostProcessor if (handlerMethodFactory != null) { this.messageHandlerMethodFactory.setMessageHandlerMethodFactory(handlerMethodFactory); } + else { + addFormatters(this.messageHandlerMethodFactory.defaultFormattingConversionService); + } // Actually register all listeners this.registrar.afterPropertiesSet(); @@ -630,6 +638,27 @@ public class KafkaListenerAnnotationBeanPostProcessor return value; } + private void addFormatters(FormatterRegistry registry) { + for (Converter converter : getBeansOfType(Converter.class)) { + registry.addConverter(converter); + } + for (GenericConverter converter : getBeansOfType(GenericConverter.class)) { + registry.addConverter(converter); + } + for (Formatter formatter : getBeansOfType(Formatter.class)) { + registry.addFormatter(formatter); + } + } + + private Collection getBeansOfType(Class type) { + if (KafkaListenerAnnotationBeanPostProcessor.this.beanFactory instanceof ListableBeanFactory) { + return ((ListableBeanFactory) KafkaListenerAnnotationBeanPostProcessor.this.beanFactory).getBeansOfType(type).values(); + } + else { + return Collections.emptySet(); + } + } + /** * An {@link MessageHandlerMethodFactory} adapter that offers a configurable underlying * instance to use. Useful if the factory to use is determined once the endpoints @@ -638,6 +667,8 @@ public class KafkaListenerAnnotationBeanPostProcessor */ private class KafkaHandlerMethodFactoryAdapter implements MessageHandlerMethodFactory { + private DefaultFormattingConversionService defaultFormattingConversionService = new DefaultFormattingConversionService(); + private MessageHandlerMethodFactory messageHandlerMethodFactory; public void setMessageHandlerMethodFactory(MessageHandlerMethodFactory kafkaHandlerMethodFactory1) { @@ -664,17 +695,17 @@ public class KafkaListenerAnnotationBeanPostProcessor (KafkaListenerAnnotationBeanPostProcessor.this.beanFactory instanceof ConfigurableBeanFactory ? (ConfigurableBeanFactory) KafkaListenerAnnotationBeanPostProcessor.this.beanFactory : null); - DefaultFormattingConversionService conversionService = new DefaultFormattingConversionService(); - defaultFactory.setConversionService(conversionService); + + defaultFactory.setConversionService(this.defaultFormattingConversionService); List argumentResolvers = new ArrayList<>(); // Annotation-based argument resolution - argumentResolvers.add(new HeaderMethodArgumentResolver(conversionService, cbf)); + argumentResolvers.add(new HeaderMethodArgumentResolver(this.defaultFormattingConversionService, cbf)); argumentResolvers.add(new HeadersMethodArgumentResolver()); // Type-based argument resolution - final GenericMessageConverter messageConverter = new GenericMessageConverter(conversionService); + final GenericMessageConverter messageConverter = new GenericMessageConverter(this.defaultFormattingConversionService); argumentResolvers.add(new MessageMethodArgumentResolver(messageConverter)); argumentResolvers.add(new PayloadArgumentResolver(messageConverter) { diff --git a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java index eca83568..ca1764c4 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2017 the original author or authors. + * Copyright 2016-2018 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. @@ -20,6 +20,8 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.anyMap; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.BDDMockito.willAnswer; +import static org.mockito.BDDMockito.willReturn; +import static org.mockito.BDDMockito.willThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; @@ -50,6 +52,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.event.EventListener; import org.springframework.context.support.PropertySourcesPlaceholderConfigurer; +import org.springframework.core.convert.converter.Converter; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.config.KafkaListenerContainerFactory; import org.springframework.kafka.config.KafkaListenerEndpointRegistry; @@ -104,6 +107,7 @@ import org.springframework.util.concurrent.ListenableFuture; * @author Artem Bilan * @author Dariusz Szablinski * @author Venil Noronha + * @author Dimitri Penner */ @ContextConfiguration @RunWith(SpringJUnit4ClassRunner.class) @@ -120,7 +124,7 @@ public class EnableKafkaIntegrationTests { "annotated18", "annotated19", "annotated20", "annotated21", "annotated21reply", "annotated22", "annotated22reply", "annotated23", "annotated23reply", "annotated24", "annotated24reply", "annotated25", "annotated25reply1", "annotated25reply2", "annotated26", "annotated27", "annotated28", - "annotated29", "annotated30", "annotated30reply", "annotated31"); + "annotated29", "annotated30", "annotated30reply", "annotated31", "annotated32"); // @Rule // public Log4jLevelAdjuster adjuster = new Log4jLevelAdjuster(Level.TRACE, @@ -159,6 +163,9 @@ public class EnableKafkaIntegrationTests { @Autowired private List quxGroup; + @Autowired + private FooConverter fooConverter; + @Test public void testAnonymous() { MessageListenerContainer container = this.registry @@ -598,6 +605,25 @@ public class EnableKafkaIntegrationTests { + "the listener container must have a MANUAL Ackmode to populate the Acknowledgment."); } + @Test + public void testConverterBean() throws Exception { + @SuppressWarnings("unchecked") + Converter converterDelegate = mock(Converter.class); + fooConverter.setDelegate(converterDelegate); + + Foo foo = new Foo(); + willReturn(foo).given(converterDelegate).convert("{'bar':'foo'}"); + template.send("annotated32", 0, 1, "{'bar':'foo'}"); + assertThat(this.listener.latch20.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(this.listener.listen16foo).isEqualTo(foo); + + willThrow(new RuntimeException()).given(converterDelegate).convert("foobar"); + template.send("annotated32", 0, 1, "foobar"); + assertThat(this.config.listen16ErrorLatch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(this.config.listen16Exception).isNotNull(); + assertThat(this.config.listen16Message).isEqualTo("foobar"); + } + @Configuration @EnableKafka @EnableTransactionManagement(proxyTargetClass = true) @@ -769,6 +795,19 @@ public class EnableKafkaIntegrationTests { return factory; } + @Bean + public KafkaListenerContainerFactory> + recordAckListenerContainerFactory() { + ConcurrentKafkaListenerContainerFactory factory = + new ConcurrentKafkaListenerContainerFactory<>(); + factory.setConsumerFactory(manualConsumerFactory("clientIdViaProps4")); + ContainerProperties props = factory.getContainerProperties(); + props.setAckMode(AckMode.RECORD); + props.setAckOnError(true); + props.setErrorHandler(listen16ErrorHandler()); + return factory; + } + @Bean public DefaultKafkaConsumerFactory consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerConfigs()); @@ -976,6 +1015,26 @@ public class EnableKafkaIntegrationTests { }; } + private Throwable listen16Exception; + + private Object listen16Message; + + private CountDownLatch listen16ErrorLatch = new CountDownLatch(1); + + @Bean + public ConsumerAwareErrorHandler listen16ErrorHandler() { + return (e, r, c) -> { + listen16Exception = e; + listen16Message = r.value(); + listen16ErrorLatch.countDown(); + }; + } + + @Bean + public FooConverter fooConverter() { + return new FooConverter(); + } + } static class Listener implements ConsumerSeekAware { @@ -1020,6 +1079,8 @@ public class EnableKafkaIntegrationTests { private final CountDownLatch latch19 = new CountDownLatch(1); + private final CountDownLatch latch20 = new CountDownLatch(1); + private final CountDownLatch eventLatch = new CountDownLatch(1); private volatile Integer partition; @@ -1036,6 +1097,8 @@ public class EnableKafkaIntegrationTests { private Foo foo; + private Foo listen16foo; + private volatile ListenerContainerIdleEvent event; private volatile List keys; @@ -1214,6 +1277,13 @@ public class EnableKafkaIntegrationTests { this.latch15.countDown(); } + @KafkaListener(id = "converter", topics = "annotated32", containerFactory = "recordAckListenerContainerFactory", + groupId = "converter.explicitGroupId") + public void listen16(Foo foo) { + this.listen16foo = foo; + this.latch20.countDown(); + } + @KafkaListener(id = "errorHandler", topics = "annotated20", errorHandler = "consumeException") public String errorHandler(String data) throws Exception { throw new Exception("return this"); @@ -1405,4 +1475,22 @@ public class EnableKafkaIntegrationTests { } + public static class FooConverter implements Converter { + + private Converter delegate; + + public Converter getDelegate() { + return delegate; + } + + public void setDelegate( + Converter delegate) { + this.delegate = delegate; + } + + @Override + public Foo convert(String source) { + return delegate.convert(source); + } + } } diff --git a/src/reference/asciidoc/kafka.adoc b/src/reference/asciidoc/kafka.adoc index f2a521f5..a0d41aff 100644 --- a/src/reference/asciidoc/kafka.adoc +++ b/src/reference/asciidoc/kafka.adoc @@ -1271,6 +1271,20 @@ public void listen(List foos, @Header(KafkaHeaders.OFFSET) List offse Notice that you can still access the batch headers too. +Starting with _versions 2.1.1_, the `org.springframework.core.convert.ConversionService` used by the default +`org.springframework.messaging.handler.annotation.support.MessageHandlerMethodFactory` to reslove parameters for the invocation +of a listener method is supplied with all beans implementing any of the following interfaces: + + - `org.springframework.core.convert.converter.Converter` + - `org.springframework.core.convert.converter.GenericConverter` + - `org.springframework.format.Formatter` + +This allows you to further customize listener deserialization without changing the default configuration for + `ConsumerFactory` and `KafkaListenerContainerFactory`. + +IMPORTANT: Setting a custom `MessageHandlerMethodFactory` on the `KafkaListenerEndpointRegistrar` through a +`KafkaListenerConfigurer` bean will disable this feature. + [[headers]] ==== Message Headers