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 d622810b5fee..9f66725b7d11 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 @@ -812,6 +812,11 @@ public static class Streams { */ private @Nullable String stateDir; + /** + * Whether the consumer should leave the group when stopping Kafka Streams. + */ + private boolean leaveGroupOnClose; + /** * Additional Kafka properties used to configure the streams. */ @@ -885,6 +890,14 @@ public void setStateDir(@Nullable String stateDir) { this.stateDir = stateDir; } + public boolean isLeaveGroupOnClose() { + return this.leaveGroupOnClose; + } + + public void setLeaveGroupOnClose(boolean leaveGroupOnClose) { + this.leaveGroupOnClose = leaveGroupOnClose; + } + public Map getProperties() { return this.properties; } diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaStreamsAnnotationDrivenConfiguration.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaStreamsAnnotationDrivenConfiguration.java index 1c7184960063..410d72067254 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaStreamsAnnotationDrivenConfiguration.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaStreamsAnnotationDrivenConfiguration.java @@ -101,6 +101,7 @@ public void configure(StreamsBuilderFactoryBean factoryBean) { KafkaProperties.Cleanup cleanup = this.properties.getStreams().getCleanup(); CleanupConfig cleanupConfig = new CleanupConfig(cleanup.isOnStartup(), cleanup.isOnShutdown()); factoryBean.setCleanupConfig(cleanupConfig); + factoryBean.setLeaveGroupOnClose(this.properties.getStreams().isLeaveGroupOnClose()); } @Override 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 e65e217272d1..f4cf275aa63b 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 @@ -475,6 +475,7 @@ void streamsBuilderFactoryBeanConfigurerIsApplied() { assertThat(factoryBean.isAutoStartup()).isFalse(); assertThat(factoryBean).extracting("cleanupConfig.onStart").isEqualTo(true); assertThat(factoryBean).extracting("cleanupConfig.onStop").isEqualTo(true); + assertThat(factoryBean).extracting("leaveGroupOnClose").isEqualTo(false); factoryBean.addListener(listener); }) .withPropertyValues("spring.kafka.client-id=cid", @@ -696,6 +697,19 @@ void streamsWithCleanupConfig() { }); } + @Test + void streamsWithLeaveGroupOnClose() { + this.contextRunner + .withUserConfiguration(EnableKafkaStreamsConfiguration.class, TestKafkaStreamsConfiguration.class) + .withPropertyValues("spring.application.name=my-test-app", + "spring.kafka.bootstrap-servers=localhost:9092,localhost:9093", + "spring.kafka.streams.auto-startup=false", "spring.kafka.streams.leave-group-on-close=true") + .run((context) -> { + StreamsBuilderFactoryBean streamsBuilderFactoryBean = context.getBean(StreamsBuilderFactoryBean.class); + assertThat(streamsBuilderFactoryBean).extracting("leaveGroupOnClose").isEqualTo(true); + }); + } + @Test void streamsApplicationIdIsMandatory() { this.contextRunner.withUserConfiguration(EnableKafkaStreamsConfiguration.class).run((context) -> {