diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/ConcurrentKafkaListenerContainerFactoryConfigurer.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/ConcurrentKafkaListenerContainerFactoryConfigurer.java index eed66b5fe3d..47b78c35726 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/ConcurrentKafkaListenerContainerFactoryConfigurer.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/ConcurrentKafkaListenerContainerFactoryConfigurer.java @@ -267,6 +267,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { map.from(properties::getLogContainerConfig).to(container::setLogContainerConfig); map.from(properties::isMissingTopicsFatal).to(container::setMissingTopicsFatal); map.from(properties::isImmediateStop).to(container::setStopImmediate); + map.from(properties::isAwaitAsyncResultsOnStop).to(container::setAwaitAsyncResultsOnStop); map.from(properties::isObservationEnabled).to(container::setObservationEnabled); map.from(properties::getAuthExceptionRetryInterval).to(container::setAuthExceptionRetryInterval); map.from(this.transactionManager).to(container::setKafkaAwareTransactionManager); diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaProperties.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaProperties.java index 9196ad5758f..9696a4ee4d7 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaProperties.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaProperties.java @@ -58,6 +58,7 @@ import org.springframework.util.unit.DataSize; * @author Andy Wilkinson * @author Scott Frederick * @author Yanming Zhou + * @author Nikita Kibitkin * @since 4.0.0 */ @ConfigurationProperties("spring.kafka") @@ -1134,6 +1135,12 @@ public class KafkaProperties { */ private boolean immediateStop; + /** + * Whether the container waits for in-flight asynchronous listener results to + * complete, within its shutdown timeout, when it stops. + */ + private boolean awaitAsyncResultsOnStop; + /** * Whether to auto start the container. */ @@ -1291,6 +1298,14 @@ public class KafkaProperties { this.immediateStop = immediateStop; } + public boolean isAwaitAsyncResultsOnStop() { + return this.awaitAsyncResultsOnStop; + } + + public void setAwaitAsyncResultsOnStop(boolean awaitAsyncResultsOnStop) { + this.awaitAsyncResultsOnStop = awaitAsyncResultsOnStop; + } + public boolean isAutoStartup() { return this.autoStartup; } diff --git a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfigurationTests.java b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfigurationTests.java index 2f1e18f1c64..82391224ddd 100644 --- a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfigurationTests.java +++ b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfigurationTests.java @@ -821,7 +821,8 @@ class KafkaAutoConfigurationTests { "spring.kafka.listener.immediate-stop=true", "spring.kafka.producer.transaction-id-prefix=foo", "spring.kafka.jaas.login-module=foo", "spring.kafka.jaas.control-flag=REQUISITE", "spring.kafka.jaas.options.useKeyTab=true", "spring.kafka.listener.async-acks=true", - "spring.kafka.listener.observation-enabled=true") + "spring.kafka.listener.observation-enabled=true", + "spring.kafka.listener.await-async-results-on-stop=true") .run((context) -> { DefaultKafkaProducerFactory producerFactory = context.getBean(DefaultKafkaProducerFactory.class); DefaultKafkaConsumerFactory consumerFactory = context.getBean(DefaultKafkaConsumerFactory.class); @@ -847,6 +848,7 @@ class KafkaAutoConfigurationTests { assertThat(containerProperties.isLogContainerConfig()).isTrue(); assertThat(containerProperties.isMissingTopicsFatal()).isTrue(); assertThat(containerProperties.isStopImmediate()).isTrue(); + assertThat(containerProperties.isAwaitAsyncResultsOnStop()).isTrue(); assertThat(containerProperties.isObservationEnabled()).isTrue(); assertThat(kafkaListenerContainerFactory).extracting("concurrency").isEqualTo(3); assertThat(kafkaListenerContainerFactory.isBatchListener()).isTrue();