From 677c05a5b144f672481d7e5591d8b41ab212e59a Mon Sep 17 00:00:00 2001 From: Michael Kreis Date: Tue, 12 Jul 2022 08:17:27 +0200 Subject: [PATCH 1/2] Add config property for KafkaAdmin modifyTopicConfigs See gh-31679 --- .../kafka/KafkaAutoConfiguration.java | 1 + .../autoconfigure/kafka/KafkaProperties.java | 13 +++++++++++++ .../kafka/KafkaAutoConfigurationTests.java | 18 +++++++++--------- 3 files changed, 23 insertions(+), 9 deletions(-) diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaAutoConfiguration.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaAutoConfiguration.java index 85329985b5..7d619429f7 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaAutoConfiguration.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaAutoConfiguration.java @@ -142,6 +142,7 @@ public class KafkaAutoConfiguration { public KafkaAdmin kafkaAdmin() { KafkaAdmin kafkaAdmin = new KafkaAdmin(this.properties.buildAdminProperties()); kafkaAdmin.setFatalIfBrokerNotAvailable(this.properties.getAdmin().isFailFast()); + kafkaAdmin.setModifyTopicConfigs(this.properties.getAdmin().isModifyTopicConfigs()); return kafkaAdmin; } diff --git a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaProperties.java b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaProperties.java index 6e48add8e0..c0db79ee5a 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaProperties.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/main/java/org/springframework/boot/autoconfigure/kafka/KafkaProperties.java @@ -647,6 +647,11 @@ public class KafkaProperties { */ private boolean failFast; + /** + * Whether to enable modification of existing topic configuration. + */ + private boolean modifyTopicConfigs; + public Ssl getSsl() { return this.ssl; } @@ -671,6 +676,14 @@ public class KafkaProperties { this.failFast = failFast; } + public boolean isModifyTopicConfigs() { + return this.modifyTopicConfigs; + } + + public void setModifyTopicConfigs(boolean modifyTopicConfigs) { + this.modifyTopicConfigs = modifyTopicConfigs; + } + public Map getProperties() { return this.properties; } 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 b61bc1528e..7682b89f2d 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 @@ -209,15 +209,14 @@ class KafkaAutoConfigurationTests { @Test void adminProperties() { - this.contextRunner - .withPropertyValues("spring.kafka.clientId=cid", "spring.kafka.properties.foo.bar.baz=qux.fiz.buz", - "spring.kafka.admin.fail-fast=true", "spring.kafka.admin.properties.fiz.buz=fix.fox", - "spring.kafka.admin.security.protocol=SSL", "spring.kafka.admin.ssl.key-password=p4", - "spring.kafka.admin.ssl.key-store-location=classpath:ksLocP", - "spring.kafka.admin.ssl.key-store-password=p5", "spring.kafka.admin.ssl.key-store-type=PKCS12", - "spring.kafka.admin.ssl.trust-store-location=classpath:tsLocP", - "spring.kafka.admin.ssl.trust-store-password=p6", - "spring.kafka.admin.ssl.trust-store-type=PKCS12", "spring.kafka.admin.ssl.protocol=TLSv1.2") + this.contextRunner.withPropertyValues("spring.kafka.clientId=cid", + "spring.kafka.properties.foo.bar.baz=qux.fiz.buz", "spring.kafka.admin.fail-fast=true", + "spring.kafka.admin.properties.fiz.buz=fix.fox", "spring.kafka.admin.security.protocol=SSL", + "spring.kafka.admin.ssl.key-password=p4", "spring.kafka.admin.ssl.key-store-location=classpath:ksLocP", + "spring.kafka.admin.ssl.key-store-password=p5", "spring.kafka.admin.ssl.key-store-type=PKCS12", + "spring.kafka.admin.ssl.trust-store-location=classpath:tsLocP", + "spring.kafka.admin.ssl.trust-store-password=p6", "spring.kafka.admin.ssl.trust-store-type=PKCS12", + "spring.kafka.admin.ssl.protocol=TLSv1.2", "spring.kafka.admin.modify-topic-configs=true") .run((context) -> { KafkaAdmin admin = context.getBean(KafkaAdmin.class); Map configs = admin.getConfigurationProperties(); @@ -239,6 +238,7 @@ class KafkaAutoConfigurationTests { assertThat(configs.get("foo.bar.baz")).isEqualTo("qux.fiz.buz"); assertThat(configs.get("fiz.buz")).isEqualTo("fix.fox"); assertThat(admin).hasFieldOrPropertyWithValue("fatalIfBrokerNotAvailable", true); + assertThat(admin).hasFieldOrPropertyWithValue("modifyTopicConfigs", true); }); } From 4ae4698093c316eee772e997ed0aef3ea37fe66f Mon Sep 17 00:00:00 2001 From: Stephane Nicoll Date: Wed, 13 Jul 2022 16:10:35 +0200 Subject: [PATCH 2/2] Polish "Add config property for KafkaAdmin modifyTopicConfigs" See gh-31679 --- .../boot/autoconfigure/kafka/KafkaPropertiesTests.java | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/kafka/KafkaPropertiesTests.java b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/kafka/KafkaPropertiesTests.java index 620e0b25ea..80a7060d77 100644 --- a/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/kafka/KafkaPropertiesTests.java +++ b/spring-boot-project/spring-boot-autoconfigure/src/test/java/org/springframework/boot/autoconfigure/kafka/KafkaPropertiesTests.java @@ -21,12 +21,14 @@ import java.util.Map; import org.apache.kafka.common.config.SslConfigs; import org.junit.jupiter.api.Test; +import org.springframework.boot.autoconfigure.kafka.KafkaProperties.Admin; import org.springframework.boot.autoconfigure.kafka.KafkaProperties.Cleanup; import org.springframework.boot.autoconfigure.kafka.KafkaProperties.IsolationLevel; import org.springframework.boot.autoconfigure.kafka.KafkaProperties.Listener; import org.springframework.boot.context.properties.source.MutuallyExclusiveConfigurationPropertiesException; import org.springframework.core.io.ClassPathResource; import org.springframework.kafka.core.CleanupConfig; +import org.springframework.kafka.core.KafkaAdmin; import org.springframework.kafka.listener.ContainerProperties; import static org.assertj.core.api.Assertions.assertThat; @@ -51,6 +53,14 @@ class KafkaPropertiesTests { assertThat(original).hasSize(IsolationLevel.values().length); } + @Test + void adminDefaultValuesAreConsistent() { + KafkaAdmin admin = new KafkaAdmin(Map.of()); + Admin adminProperties = new KafkaProperties().getAdmin(); + assertThat(admin).hasFieldOrPropertyWithValue("fatalIfBrokerNotAvailable", adminProperties.isFailFast()); + assertThat(admin).hasFieldOrPropertyWithValue("modifyTopicConfigs", adminProperties.isModifyTopicConfigs()); + } + @Test void listenerDefaultValuesAreConsistent() { ContainerProperties container = new ContainerProperties("test");