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 1631f5980d..78d614f999 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 @@ -25,6 +25,7 @@ import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.AfterRollbackProcessor; import org.springframework.kafka.listener.BatchErrorHandler; +import org.springframework.kafka.listener.ConsumerAwareRebalanceListener; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.listener.ErrorHandler; import org.springframework.kafka.support.converter.MessageConverter; @@ -47,6 +48,8 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { private KafkaAwareTransactionManager transactionManager; + private ConsumerAwareRebalanceListener rebalanceListener; + private ErrorHandler errorHandler; private BatchErrorHandler batchErrorHandler; @@ -86,6 +89,15 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { this.transactionManager = transactionManager; } + /** + * Set the {@link ConsumerAwareRebalanceListener} to use. + * @param rebalanceListener the rebalance listener. + * @since 2.2 + */ + void setRebalanceListener(ConsumerAwareRebalanceListener rebalanceListener) { + this.rebalanceListener = rebalanceListener; + } + /** * Set the {@link ErrorHandler} to use. * @param errorHandler the error handler @@ -160,6 +172,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { map.from(properties::getLogContainerConfig).to(container::setLogContainerConfig); map.from(properties::isMissingTopicsFatal).to(container::setMissingTopicsFatal); map.from(this.transactionManager).to(container::setTransactionManager); + map.from(this.rebalanceListener).to(container::setConsumerRebalanceListener); } } 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 34e78548c9..89dca80442 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 @@ -29,6 +29,7 @@ import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.AfterRollbackProcessor; import org.springframework.kafka.listener.BatchErrorHandler; +import org.springframework.kafka.listener.ConsumerAwareRebalanceListener; import org.springframework.kafka.listener.ErrorHandler; import org.springframework.kafka.support.converter.BatchMessageConverter; import org.springframework.kafka.support.converter.BatchMessagingMessageConverter; @@ -57,6 +58,8 @@ class KafkaAnnotationDrivenConfiguration { private final KafkaAwareTransactionManager transactionManager; + private final ConsumerAwareRebalanceListener rebalanceListener; + private final ErrorHandler errorHandler; private final BatchErrorHandler batchErrorHandler; @@ -68,6 +71,7 @@ class KafkaAnnotationDrivenConfiguration { ObjectProvider batchMessageConverter, ObjectProvider> kafkaTemplate, ObjectProvider> kafkaTransactionManager, + ObjectProvider rebalanceListener, ObjectProvider errorHandler, ObjectProvider batchErrorHandler, ObjectProvider> afterRollbackProcessor) { @@ -77,6 +81,7 @@ class KafkaAnnotationDrivenConfiguration { () -> new BatchMessagingMessageConverter(this.messageConverter)); this.kafkaTemplate = kafkaTemplate.getIfUnique(); this.transactionManager = kafkaTransactionManager.getIfUnique(); + this.rebalanceListener = rebalanceListener.getIfUnique(); this.errorHandler = errorHandler.getIfUnique(); this.batchErrorHandler = batchErrorHandler.getIfUnique(); this.afterRollbackProcessor = afterRollbackProcessor.getIfUnique(); @@ -92,6 +97,7 @@ class KafkaAnnotationDrivenConfiguration { configurer.setMessageConverter(messageConverterToUse); configurer.setReplyTemplate(this.kafkaTemplate); configurer.setTransactionManager(this.transactionManager); + configurer.setRebalanceListener(this.rebalanceListener); configurer.setErrorHandler(this.errorHandler); configurer.setBatchErrorHandler(this.batchErrorHandler); configurer.setAfterRollbackProcessor(this.afterRollbackProcessor); 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 935af69698..5c1c385904 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 @@ -55,6 +55,7 @@ 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.ConsumerAwareRebalanceListener; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.listener.ContainerProperties.AckMode; import org.springframework.kafka.listener.SeekToCurrentBatchErrorHandler; @@ -674,6 +675,18 @@ public class KafkaAutoConfigurationTests { }); } + @Test + public void testConcurrentKafkaListenerContainerFactoryWithCustomRebalanceListener() { + this.contextRunner.withUserConfiguration(RebalanceListenerConfiguration.class) + .run((context) -> { + ConcurrentKafkaListenerContainerFactory factory = context + .getBean(ConcurrentKafkaListenerContainerFactory.class); + assertThat(factory.getContainerProperties()) + .hasFieldOrPropertyWithValue("consumerRebalanceListener", + context.getBean("rebalanceListener")); + }); + } + @Test public void testConcurrentKafkaListenerContainerFactoryWithKafkaTemplate() { this.contextRunner.run((context) -> { @@ -749,6 +762,16 @@ public class KafkaAutoConfigurationTests { } + @Configuration(proxyBeanMethods = false) + protected static class RebalanceListenerConfiguration { + + @Bean + public ConsumerAwareRebalanceListener rebalanceListener() { + return mock(ConsumerAwareRebalanceListener.class); + } + + } + @Configuration(proxyBeanMethods = false) @EnableKafkaStreams protected static class EnableKafkaStreamsConfiguration { 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 15550b8e3c..728d4a4caa 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 @@ -6158,8 +6158,9 @@ The following component creates a listener endpoint on the `someTopic` topic: ---- If a `KafkaTransactionManager` bean is defined, it is automatically associated to the -container factory. Similarly, if a `ErrorHandler` or `AfterRollbackProcessor` bean is -defined, it is automatically associated to the default factory. +container factory. Similarly, if a `ErrorHandler`, `AfterRollbackProcessor` or +`ConsumerAwareRebalanceListener` bean is defined, it is automatically associated to the +default factory. Depending on the listener type, a `RecordMessageConverter` or `BatchMessageConverter` bean is associated to the default factory. If only a `RecordMessageConverter` bean is present