From 3c6f7a54d30e82c9a06b635bb79d9a4624fa6633 Mon Sep 17 00:00:00 2001 From: Nikita Kibitkin Date: Tue, 15 Sep 2026 19:47:16 +0200 Subject: [PATCH] Add property to await async listener results on stop Spring Kafka 4.2 adds ContainerProperties#setAwaitAsyncResultsOnStop: when the container stops, in-flight asynchronous listener results (CompletableFuture, Mono, Kotlin suspend functions) are awaited within the shutdown timeout and cancelled afterwards, so that listener work does not outlive the container. This commit adds spring.kafka.listener.await-async-results-on-stop so that this can be enabled with configuration rather than a ContainerCustomizer. See gh-51774 Signed-off-by: Nikita Kibitkin --- ...ntKafkaListenerContainerFactoryConfigurer.java | 1 + .../boot/kafka/autoconfigure/KafkaProperties.java | 15 +++++++++++++++ .../KafkaAutoConfigurationTests.java | 4 +++- 3 files changed, 19 insertions(+), 1 deletion(-) 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();