From 59c6dc5c7ad595082ccf37ddd8f891df86543b21 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 27 Aug 2018 09:27:11 -0400 Subject: [PATCH] Improve Kafka Auto-configuration - transaction manager - error handler - after rollback processor See gh-14215 --- ...fkaListenerContainerFactoryConfigurer.java | 40 +++++++++++++++++++ .../KafkaAnnotationDrivenConfiguration.java | 20 +++++++++- .../kafka/KafkaAutoConfigurationTests.java | 37 +++++++++++++++++ .../spring-boot-dependencies/pom.xml | 2 +- .../main/asciidoc/spring-boot-features.adoc | 15 ++++++- 5 files changed, 110 insertions(+), 4 deletions(-) diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/ConcurrentKafkaListenerContainerFactoryConfigurer.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/ConcurrentKafkaListenerContainerFactoryConfigurer.java index bdc00693f5..c6b3827351 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/ConcurrentKafkaListenerContainerFactoryConfigurer.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/ConcurrentKafkaListenerContainerFactoryConfigurer.java @@ -23,8 +23,11 @@ import org.springframework.boot.context.properties.PropertyMapper; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.listener.AfterRollbackProcessor; import org.springframework.kafka.listener.ContainerProperties; +import org.springframework.kafka.listener.ErrorHandler; import org.springframework.kafka.support.converter.RecordMessageConverter; +import org.springframework.kafka.transaction.KafkaAwareTransactionManager; /** * Configure {@link ConcurrentKafkaListenerContainerFactory} with sensible defaults. @@ -41,6 +44,12 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { private KafkaTemplate replyTemplate; + private KafkaAwareTransactionManager transactionManager; + + private ErrorHandler errorHandler; + + private AfterRollbackProcessor afterRollbackProcessor; + /** * Set the {@link KafkaProperties} to use. * @param properties the properties @@ -65,6 +74,32 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { this.replyTemplate = replyTemplate; } + /** + * Set the {@link KafkaAwareTransactionManager} to use. + * @param transactionManager the transaction manager + */ + public void setTransactionManager( + KafkaAwareTransactionManager transactionManager) { + this.transactionManager = transactionManager; + } + + /** + * Set the {@link ErrorHandler} to use. + * @param errorHandler the error handler + */ + public void setErrorHandler(ErrorHandler errorHandler) { + this.errorHandler = errorHandler; + } + + /** + * Set the {@link AfterRollbackProcessor} to use. + * @param afterRollbackProcessor the after rollback processor + */ + public void setAfterRollbackProcessor( + AfterRollbackProcessor afterRollbackProcessor) { + this.afterRollbackProcessor = afterRollbackProcessor; + } + /** * Configure the specified Kafka listener container factory. The factory can be * further tuned and default settings can be overridden. @@ -89,6 +124,9 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { map.from(this.replyTemplate).whenNonNull().to(factory::setReplyTemplate); map.from(properties::getType).whenEqualTo(Listener.Type.BATCH) .toCall(() -> factory.setBatchListener(true)); + map.from(this.errorHandler).whenNonNull().to(factory::setErrorHandler); + map.from(this.afterRollbackProcessor).whenNonNull() + .to(factory::setAfterRollbackProcessor); } private void configureContainer(ContainerProperties container) { @@ -109,6 +147,8 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { .as(Number::intValue).to(container::setMonitorInterval); map.from(properties::getLogContainerConfig).whenNonNull() .to(container::setLogContainerConfig); + map.from(this.transactionManager).whenNonNull() + .to(container::setTransactionManager); } } diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaAnnotationDrivenConfiguration.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaAnnotationDrivenConfiguration.java index 57ee0c9c7e..a9b74a0332 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaAnnotationDrivenConfiguration.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaAnnotationDrivenConfiguration.java @@ -26,7 +26,10 @@ import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.config.KafkaListenerConfigUtils; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.listener.AfterRollbackProcessor; +import org.springframework.kafka.listener.ErrorHandler; import org.springframework.kafka.support.converter.RecordMessageConverter; +import org.springframework.kafka.transaction.KafkaAwareTransactionManager; /** * Configuration for Kafka annotation-driven support. @@ -45,12 +48,24 @@ class KafkaAnnotationDrivenConfiguration { private final KafkaTemplate kafkaTemplate; + private final KafkaAwareTransactionManager transactionManager; + + private final ErrorHandler errorHandler; + + private final AfterRollbackProcessor afterRollbackProcessor; + KafkaAnnotationDrivenConfiguration(KafkaProperties properties, ObjectProvider messageConverter, - ObjectProvider> kafkaTemplate) { + ObjectProvider> kafkaTemplate, + ObjectProvider> kafkaTransactionManager, + ObjectProvider errorHandler, + ObjectProvider> afterRollbackProcessor) { this.properties = properties; this.messageConverter = messageConverter.getIfUnique(); this.kafkaTemplate = kafkaTemplate.getIfUnique(); + this.transactionManager = kafkaTransactionManager.getIfUnique(); + this.errorHandler = errorHandler.getIfUnique(); + this.afterRollbackProcessor = afterRollbackProcessor.getIfUnique(); } @Bean @@ -60,6 +75,9 @@ class KafkaAnnotationDrivenConfiguration { configurer.setKafkaProperties(this.properties); configurer.setMessageConverter(this.messageConverter); configurer.setReplyTemplate(this.kafkaTemplate); + configurer.setErrorHandler(this.errorHandler); + configurer.setTransactionManager(this.transactionManager); + configurer.setAfterRollbackProcessor(this.afterRollbackProcessor); return configurer; } diff --git a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/kafka/KafkaAutoConfigurationTests.java b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/kafka/KafkaAutoConfigurationTests.java index 816559e886..8c168af3c6 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/kafka/KafkaAutoConfigurationTests.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/kafka/KafkaAutoConfigurationTests.java @@ -41,6 +41,7 @@ import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.test.context.runner.ApplicationContextRunner; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Primary; import org.springframework.kafka.annotation.EnableKafkaStreams; import org.springframework.kafka.annotation.KafkaStreamsDefaultConfiguration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; @@ -50,12 +51,16 @@ import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaAdmin; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.listener.AfterRollbackProcessor; import org.springframework.kafka.listener.ContainerProperties.AckMode; +import org.springframework.kafka.listener.SeekToCurrentErrorHandler; import org.springframework.kafka.security.jaas.KafkaJaasLoginModuleInitializer; import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.kafka.test.utils.KafkaTestUtils; +import org.springframework.kafka.transaction.ChainedKafkaTransactionManager; import org.springframework.kafka.transaction.KafkaTransactionManager; +import org.springframework.transaction.PlatformTransactionManager; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.entry; @@ -155,6 +160,10 @@ public class KafkaAutoConfigurationTests { assertThat(configs.get("baz")).isEqualTo("qux"); assertThat(configs.get("foo.bar.baz")).isEqualTo("qux.fiz.buz"); assertThat(configs.get("fiz.buz")).isEqualTo("fix.fox"); + ConcurrentKafkaListenerContainerFactory factory = context + .getBean(ConcurrentKafkaListenerContainerFactory.class); + assertThat(KafkaTestUtils.getPropertyValue(factory, "errorHandler")) + .isSameAs(context.getBean("errorHandler")); }); } @@ -485,6 +494,7 @@ public class KafkaAutoConfigurationTests { @Test public void testKafkaTemplateRecordMessageConverters() { this.contextRunner.withUserConfiguration(MessageConverterConfiguration.class) + .withPropertyValues("spring.kafka.producer.transaction-id-prefix=test") .run((context) -> { KafkaTemplate kafkaTemplate = context .getBean(KafkaTemplate.class); @@ -496,6 +506,7 @@ public class KafkaAutoConfigurationTests { @Test public void testConcurrentKafkaListenerContainerFactoryWithCustomMessageConverters() { this.contextRunner.withUserConfiguration(MessageConverterConfiguration.class) + .withPropertyValues("spring.kafka.producer.transaction-id-prefix=test") .run((context) -> { ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory = context .getBean(ConcurrentKafkaListenerContainerFactory.class); @@ -503,6 +514,11 @@ public class KafkaAutoConfigurationTests { kafkaListenerContainerFactory); assertThat(dfa.getPropertyValue("messageConverter")) .isSameAs(context.getBean("myMessageConverter")); + assertThat(kafkaListenerContainerFactory.getContainerProperties() + .getTransactionManager()).isSameAs( + context.getBean("chainedTransactionManager")); + assertThat(dfa.getPropertyValue("afterRollbackProcessor")) + .isSameAs(context.getBean("arp")); }); } @@ -521,6 +537,11 @@ public class KafkaAutoConfigurationTests { @Configuration protected static class TestConfiguration { + @Bean + public SeekToCurrentErrorHandler errorHandler() { + return new SeekToCurrentErrorHandler(); + } + } @Configuration @@ -531,6 +552,22 @@ public class KafkaAutoConfigurationTests { return mock(RecordMessageConverter.class); } + @Bean + @Primary + public PlatformTransactionManager chainedTransactionManager( + KafkaTransactionManager kafkaTransactionManager) { + + return new ChainedKafkaTransactionManager( + kafkaTransactionManager); + } + + @Bean + public AfterRollbackProcessor arp() { + return (records, consumer, ex, recoverable) -> { + // no-op + }; + } + } @Configuration diff --git a/spring-boot-project/spring-boot-dependencies/pom.xml b/spring-boot-project/spring-boot-dependencies/pom.xml index 8ccaae1407..8868b75aaa 100644 --- a/spring-boot-project/spring-boot-dependencies/pom.xml +++ b/spring-boot-project/spring-boot-dependencies/pom.xml @@ -161,7 +161,7 @@ Lovelace-RC2 0.25.0.RELEASE 5.1.0.M2 - 2.2.0.M2 + 2.2.0.BUILD-SNAPSHOT 2.3.2.RELEASE 1.2.0.RELEASE 2.0.2.RELEASE diff --git a/spring-boot-project/spring-boot-docs/src/main/asciidoc/spring-boot-features.adoc b/spring-boot-project/spring-boot-docs/src/main/asciidoc/spring-boot-features.adoc index d871f00702..8219e82528 100644 --- a/spring-boot-project/spring-boot-docs/src/main/asciidoc/spring-boot-features.adoc +++ b/spring-boot-project/spring-boot-docs/src/main/asciidoc/spring-boot-features.adoc @@ -5627,8 +5627,19 @@ bean is defined, it is automatically associated to the auto-configured `KafkaTem When the Apache Kafka infrastructure is present, any bean can be annotated with `@KafkaListener` to create a listener endpoint. If no `KafkaListenerContainerFactory` has been defined, a default one is automatically configured with keys defined in -`spring.kafka.listener.*`. Also, if a `RecordMessageConverter` bean is defined, it is -automatically associated to the default factory. +`spring.kafka.listener.*`. +If the property `spring.kafka.producer.transaction-id-prefix` is defined, a +`KafkaTransactionManager` will be auto-configured with name `kafkaTransactionManager` and +will be wired into the container factory. +Also, if `RecordMessageConverter`, `ErrorHandler` and/or +`AfterRollbackProcessor` beans are defined, they are automatically associated to the +default factory. + +IMPORTANT: The auto configuration of these beans occur if there is just a single +instance. +When using a `ChainedKafkaTransactionManager`, it will usually reference the configured +`KafkaTransactionManager` bean, so the chained manager must be marked +`@Primary` if you want it wired into the container factory. The following component creates a listener endpoint on the `someTopic` topic: