From dcbb307d3e879f0c4bce37e68888053a75ef7e4a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?St=C3=A9phane=20Nicoll?= Date: Tue, 15 Sep 2026 14:58:01 +0200 Subject: [PATCH] Add support for dedicated KafkaAdmin for listener and template This commit adds a spring.kafka.listener.admin and spring.kafka.template.admin namespace with similar configuration properties than spring.kafka.admin, used to create the KafkaAdmin bean. When a property is configured in those two namespaces, a dedicated KafkaAdmin is created and associated with the relevant component. Closes gh-38830 --- ...fkaListenerContainerFactoryConfigurer.java | 12 +++ .../autoconfigure/KafkaAdminBuilder.java | 57 +++++++++++++ .../KafkaAnnotationDrivenConfiguration.java | 14 +++- .../autoconfigure/KafkaAutoConfiguration.java | 29 +++---- .../autoconfigure/KafkaConfigBuilder.java | 17 +++- .../kafka/autoconfigure/KafkaProperties.java | 80 ++++++++++++------- ...stenerContainerFactoryConfigurerTests.java | 10 +++ .../KafkaAutoConfigurationTests.java | 39 +++++++++ .../KafkaConfigBuilderTests.java | 12 +++ 9 files changed, 221 insertions(+), 49 deletions(-) create mode 100644 module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAdminBuilder.java 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 df90651c947..eed66b5fe3d 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 @@ -26,6 +26,7 @@ import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Listener; import org.springframework.core.task.SimpleAsyncTaskExecutor; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; +import org.springframework.kafka.core.KafkaAdmin; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.AfterRollbackProcessor; import org.springframework.kafka.listener.BatchInterceptor; @@ -59,6 +60,8 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { private @Nullable KafkaProperties properties; + private @Nullable KafkaAdmin kafkaAdmin; + private @Nullable BatchMessageConverter batchMessageConverter; private @Nullable RecordMessageConverter recordMessageConverter; @@ -93,6 +96,14 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { this.properties = properties; } + /** + * Set the {@link KafkaAdmin} to use. + * @param kafkaAdmin the Kafka admin + */ + void setKafkaAdmin(@Nullable KafkaAdmin kafkaAdmin) { + this.kafkaAdmin = kafkaAdmin; + } + /** * Set the {@link BatchMessageConverter} to use. * @param batchMessageConverter the message converter @@ -217,6 +228,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { Listener properties = this.properties.getListener(); map.from(properties::getConcurrency).to(factory::setConcurrency); map.from(properties::isAutoStartup).to(factory::setAutoStartup); + map.from(this.kafkaAdmin).to(factory::setKafkaAdmin); map.from(this.batchMessageConverter).to(factory::setBatchMessageConverter); map.from(this.recordMessageConverter).to(factory::setRecordMessageConverter); map.from(this.recordFilterStrategy).to(factory::setRecordFilterStrategy); diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAdminBuilder.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAdminBuilder.java new file mode 100644 index 00000000000..78bc264240e --- /dev/null +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAdminBuilder.java @@ -0,0 +1,57 @@ +/* + * Copyright 2012-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.boot.kafka.autoconfigure; + +import java.util.Map; + +import org.jspecify.annotations.Nullable; + +import org.springframework.boot.kafka.autoconfigure.KafkaProperties.SimpleAdmin; +import org.springframework.kafka.core.KafkaAdmin; + +/** + * Builder for {@link KafkaAdmin}. + * + * @author Stephane Nicoll + */ +class KafkaAdminBuilder { + + private final KafkaConfigBuilder kafkaConfigBuilder; + + private final @Nullable KafkaConnectionDetails connectionDetails; + + KafkaAdminBuilder(KafkaProperties properties, @Nullable KafkaConnectionDetails connectionDetails) { + this.kafkaConfigBuilder = KafkaConfigBuilder.of(properties); + this.connectionDetails = connectionDetails; + } + + KafkaAdmin build(SimpleAdmin adminProperties) { + Map properties = this.kafkaConfigBuilder.admin(adminProperties) + .withConnectionDetails(this.connectionDetails) + .build(); + KafkaAdmin kafkaAdmin = new KafkaAdmin(properties); + if (adminProperties.getCloseTimeout() != null) { + kafkaAdmin.setCloseTimeout((int) adminProperties.getCloseTimeout().getSeconds()); + } + if (adminProperties.getOperationTimeout() != null) { + kafkaAdmin.setOperationTimeout((int) adminProperties.getOperationTimeout().getSeconds()); + } + kafkaAdmin.setFatalIfBrokerNotAvailable(adminProperties.isFailFast()); + return kafkaAdmin; + } + +} diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAnnotationDrivenConfiguration.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAnnotationDrivenConfiguration.java index ffecc542e62..eb2f7ce3cd8 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAnnotationDrivenConfiguration.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAnnotationDrivenConfiguration.java @@ -24,6 +24,7 @@ import org.springframework.beans.factory.ObjectProvider; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnThreading; +import org.springframework.boot.kafka.autoconfigure.KafkaProperties.SimpleAdmin; import org.springframework.boot.thread.Threading; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -34,6 +35,7 @@ import org.springframework.kafka.config.ContainerCustomizer; import org.springframework.kafka.config.KafkaListenerConfigUtils; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; +import org.springframework.kafka.core.KafkaAdmin; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.AfterRollbackProcessor; import org.springframework.kafka.listener.BatchInterceptor; @@ -69,6 +71,8 @@ class KafkaAnnotationDrivenConfiguration { private final @Nullable RecordFilterStrategy recordFilterStrategy; + private final KafkaConnectionDetails kafkaConnectionDetails; + private final BatchMessageConverter batchMessageConverter; private final @Nullable KafkaTemplate kafkaTemplate; @@ -92,7 +96,7 @@ class KafkaAnnotationDrivenConfiguration { KafkaAnnotationDrivenConfiguration(KafkaProperties properties, ObjectProvider recordMessageConverter, ObjectProvider> recordFilterStrategy, - ObjectProvider batchMessageConverter, + KafkaConnectionDetails kafkaConnectionDetails, ObjectProvider batchMessageConverter, ObjectProvider> kafkaTemplate, ObjectProvider> kafkaTransactionManager, ObjectProvider rebalanceListener, @@ -107,6 +111,7 @@ class KafkaAnnotationDrivenConfiguration { this.recordFilterStrategy = recordFilterStrategy.getIfUnique(); this.batchMessageConverter = batchMessageConverter .getIfUnique(() -> new BatchMessagingMessageConverter(this.recordMessageConverter)); + this.kafkaConnectionDetails = kafkaConnectionDetails; this.kafkaTemplate = kafkaTemplate.getIfUnique(); this.transactionManager = kafkaTransactionManager.getIfUnique(); this.rebalanceListener = rebalanceListener.getIfUnique(); @@ -139,6 +144,7 @@ class KafkaAnnotationDrivenConfiguration { private ConcurrentKafkaListenerContainerFactoryConfigurer configurer() { ConcurrentKafkaListenerContainerFactoryConfigurer configurer = new ConcurrentKafkaListenerContainerFactoryConfigurer(); configurer.setKafkaProperties(this.properties); + configurer.setKafkaAdmin(createKafkaAdmin()); configurer.setBatchMessageConverter(this.batchMessageConverter); configurer.setRecordMessageConverter(this.recordMessageConverter); configurer.setRecordFilterStrategy(this.recordFilterStrategy); @@ -154,6 +160,12 @@ class KafkaAnnotationDrivenConfiguration { return configurer; } + private @Nullable KafkaAdmin createKafkaAdmin() { + SimpleAdmin adminProperties = this.properties.getListener().getAdmin(); + return (adminProperties != null) + ? new KafkaAdminBuilder(this.properties, this.kafkaConnectionDetails).build(adminProperties) : null; + } + @Bean @ConditionalOnMissingBean(name = "kafkaListenerContainerFactory") ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory( diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfiguration.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfiguration.java index 41f359e4f48..96c0a17adaa 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfiguration.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfiguration.java @@ -36,8 +36,10 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.autoconfigure.condition.ConditionalOnSingleCandidate; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.boot.context.properties.PropertyMapper; +import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Admin; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Jaas; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Retry.Topic.Backoff; +import org.springframework.boot.kafka.autoconfigure.KafkaProperties.SimpleAdmin; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Template; import org.springframework.boot.ssl.SslBundles; import org.springframework.context.annotation.Bean; @@ -99,7 +101,8 @@ public final class KafkaAutoConfiguration { KafkaTemplate kafkaTemplate(ProducerFactory kafkaProducerFactory, ProducerListener kafkaProducerListener, ObjectProvider messageConverter, - ObjectProvider observationConvention) { + ObjectProvider observationConvention, + KafkaConnectionDetails kafkaConnectionDetails) { PropertyMapper map = PropertyMapper.get(); KafkaTemplate kafkaTemplate = new KafkaTemplate<>(kafkaProducerFactory); messageConverter.ifUnique(kafkaTemplate::setMessageConverter); @@ -111,6 +114,11 @@ public final class KafkaAutoConfiguration { map.from(templateProperties.getCloseTimeout()).to(kafkaTemplate::setCloseTimeout); map.from(templateProperties.isAllowNonTransactional()).to(kafkaTemplate::setAllowNonTransactional); map.from(templateProperties.isObservationEnabled()).to(kafkaTemplate::setObservationEnabled); + SimpleAdmin adminProperties = templateProperties.getAdmin(); + if (adminProperties != null) { + kafkaTemplate + .setKafkaAdmin(new KafkaAdminBuilder(this.properties, kafkaConnectionDetails).build(adminProperties)); + } return kafkaTemplate; } @@ -176,21 +184,10 @@ public final class KafkaAutoConfiguration { @Bean @ConditionalOnMissingBean KafkaAdmin kafkaAdmin(KafkaConnectionDetails connectionDetails) { - Map properties = KafkaConfigBuilder.of(this.properties) - .admin() - .withConnectionDetails(connectionDetails) - .build(); - KafkaAdmin kafkaAdmin = new KafkaAdmin(properties); - KafkaProperties.Admin admin = this.properties.getAdmin(); - if (admin.getCloseTimeout() != null) { - kafkaAdmin.setCloseTimeout((int) admin.getCloseTimeout().getSeconds()); - } - if (admin.getOperationTimeout() != null) { - kafkaAdmin.setOperationTimeout((int) admin.getOperationTimeout().getSeconds()); - } - kafkaAdmin.setFatalIfBrokerNotAvailable(admin.isFailFast()); - kafkaAdmin.setModifyTopicConfigs(admin.isModifyTopicConfigs()); - kafkaAdmin.setAutoCreate(admin.isAutoCreate()); + Admin adminProperties = this.properties.getAdmin(); + KafkaAdmin kafkaAdmin = new KafkaAdminBuilder(this.properties, connectionDetails).build(adminProperties); + kafkaAdmin.setModifyTopicConfigs(adminProperties.isModifyTopicConfigs()); + kafkaAdmin.setAutoCreate(adminProperties.isAutoCreate()); return kafkaAdmin; } diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilder.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilder.java index de1fd4bdf6c..d6e7de664d5 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilder.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilder.java @@ -35,7 +35,6 @@ import org.jspecify.annotations.Nullable; import org.springframework.boot.context.properties.PropertyMapper; import org.springframework.boot.context.properties.source.MutuallyExclusiveConfigurationPropertiesException; import org.springframework.boot.kafka.autoconfigure.KafkaConnectionDetails.Configuration; -import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Admin; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Producer; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Security; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Ssl; @@ -76,7 +75,17 @@ public class KafkaConfigBuilder { * @see AdminClientConfig */ public ConfigBuilder admin() { - return new AdminConfigBuilder(initializeKafkaConfig(), this.kafkaProperties.getAdmin(), null); + return admin(this.kafkaProperties.getAdmin()); + } + + /** + * Return a builder for Admin-related configuration. + * @param adminProperties the admin properties to map + * @return an admin config builder + * @see AdminClientConfig + */ + public ConfigBuilder admin(KafkaProperties.SimpleAdmin adminProperties) { + return new AdminConfigBuilder(initializeKafkaConfig(), adminProperties, null); } /** @@ -175,11 +184,11 @@ public class KafkaConfigBuilder { private final KafkaConfig kafkaConfig; - private final KafkaProperties.Admin admin; + private final KafkaProperties.SimpleAdmin admin; private final @Nullable KafkaConnectionDetails connectionDetails; - private AdminConfigBuilder(KafkaConfig kafkaConfig, Admin admin, + private AdminConfigBuilder(KafkaConfig kafkaConfig, KafkaProperties.SimpleAdmin admin, @Nullable KafkaConnectionDetails connectionDetails) { this.kafkaConfig = kafkaConfig; this.admin = admin; 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 8ae3d38042b..9196ad5758f 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 @@ -663,7 +663,7 @@ public class KafkaProperties { } - public static class Admin { + public static class SimpleAdmin { private final Ssl ssl = new Ssl(); @@ -694,17 +694,6 @@ public class KafkaProperties { */ private boolean failFast; - /** - * Whether to enable modification of existing topic configuration. - */ - private boolean modifyTopicConfigs; - - /** - * Whether to automatically create topics during context initialization. When set - * to false, disables automatic topic creation during context initialization. - */ - private boolean autoCreate = true; - public Ssl getSsl() { return this.ssl; } @@ -745,22 +734,6 @@ public class KafkaProperties { this.failFast = failFast; } - public boolean isModifyTopicConfigs() { - return this.modifyTopicConfigs; - } - - public void setModifyTopicConfigs(boolean modifyTopicConfigs) { - this.modifyTopicConfigs = modifyTopicConfigs; - } - - public boolean isAutoCreate() { - return this.autoCreate; - } - - public void setAutoCreate(boolean autoCreate) { - this.autoCreate = autoCreate; - } - public Map getProperties() { return this.properties; } @@ -782,6 +755,37 @@ public class KafkaProperties { } + public static class Admin extends SimpleAdmin { + + /** + * Whether to enable modification of existing topic configuration. + */ + private boolean modifyTopicConfigs; + + /** + * Whether to automatically create topics during context initialization. When set + * to false, disables automatic topic creation during context initialization. + */ + private boolean autoCreate = true; + + public boolean isModifyTopicConfigs() { + return this.modifyTopicConfigs; + } + + public void setModifyTopicConfigs(boolean modifyTopicConfigs) { + this.modifyTopicConfigs = modifyTopicConfigs; + } + + public boolean isAutoCreate() { + return this.autoCreate; + } + + public void setAutoCreate(boolean autoCreate) { + this.autoCreate = autoCreate; + } + + } + /** * High (and some medium) priority Streams properties and a general properties bucket. */ @@ -946,6 +950,8 @@ public class KafkaProperties { public static class Template { + private @Nullable SimpleAdmin admin; + /** * Default topic to which messages are sent. */ @@ -972,6 +978,14 @@ public class KafkaProperties { */ private boolean observationEnabled; + public @Nullable SimpleAdmin getAdmin() { + return this.admin; + } + + public void setAdmin(@Nullable SimpleAdmin admin) { + this.admin = admin; + } + public @Nullable String getDefaultTopic() { return this.defaultTopic; } @@ -1030,6 +1044,8 @@ public class KafkaProperties { } + private @Nullable SimpleAdmin admin; + /** * Listener type. */ @@ -1139,6 +1155,14 @@ public class KafkaProperties { */ private @Nullable Duration authExceptionRetryInterval; + public @Nullable SimpleAdmin getAdmin() { + return this.admin; + } + + public void setAdmin(@Nullable SimpleAdmin admin) { + this.admin = admin; + } + public Type getType() { return this.type; } diff --git a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/ConcurrentKafkaListenerContainerFactoryConfigurerTests.java b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/ConcurrentKafkaListenerContainerFactoryConfigurerTests.java index 932ee15a73d..33d847a0125 100644 --- a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/ConcurrentKafkaListenerContainerFactoryConfigurerTests.java +++ b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/ConcurrentKafkaListenerContainerFactoryConfigurerTests.java @@ -25,6 +25,8 @@ import org.junit.jupiter.api.Test; import org.springframework.core.task.SimpleAsyncTaskExecutor; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; +import org.springframework.kafka.core.KafkaAdmin; +import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.MessageListenerContainer; import org.springframework.kafka.support.micrometer.KafkaListenerObservation.DefaultKafkaListenerObservationConvention; @@ -56,7 +58,15 @@ class ConcurrentKafkaListenerContainerFactoryConfigurerTests { this.configurer.setKafkaProperties(this.properties); this.factory = spy(new ConcurrentKafkaListenerContainerFactory<>()); this.consumerFactory = mock(ConsumerFactory.class); + } + @Test + void shouldApplyKafkaAdmin() { + KafkaAdmin kafkaAdmin = mock(); + this.configurer.setKafkaAdmin(kafkaAdmin); + this.configurer.configure(this.factory, this.consumerFactory); + ConcurrentMessageListenerContainer container = this.factory.createContainer("test"); + assertThat(container.getKafkaAdmin()).isSameAs(kafkaAdmin); } @Test 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 7e8bf4cf386..2f1e18f1c64 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 @@ -779,6 +779,7 @@ class KafkaAutoConfigurationTests { "spring.kafka.template.allow-non-transactional=true", "spring.kafka.template.observation-enabled=true") .run((context) -> { KafkaTemplate kafkaTemplate = context.getBean(KafkaTemplate.class); + assertThat(kafkaTemplate.getKafkaAdmin()).isSameAs(context.getBean(KafkaAdmin.class)); assertThat(kafkaTemplate.getDefaultTopic()).isEqualTo("testTopic"); assertThat(kafkaTemplate).hasFieldOrPropertyWithValue("transactionIdPrefix", "txOverride"); assertThat(kafkaTemplate).hasFieldOrPropertyWithValue("closeTimeout", Duration.ofMinutes(3)); @@ -787,6 +788,24 @@ class KafkaAutoConfigurationTests { }); } + @Test + void templatePropertiesWithAdmin() { + this.contextRunner + .withPropertyValues("spring.kafka.template.admin.client-id=template", + "spring.kafka.template.admin.operation-timeout=2m", + "spring.kafka.template.admin.properties.fiz.buz=fix.fox") + .run((context) -> { + KafkaTemplate kafkaTemplate = context.getBean(KafkaTemplate.class); + assertThat(kafkaTemplate.getKafkaAdmin()).satisfies((kafkaAdmin) -> { + assertThat(kafkaAdmin).isNotSameAs(context.getBean(KafkaAdmin.class)); + assertThat(kafkaAdmin.getConfigurationProperties()) + .containsEntry(AdminClientConfig.CLIENT_ID_CONFIG, "template") + .containsEntry("fiz.buz", "fix.fox"); + assertThat(kafkaAdmin.getOperationTimeout()).isEqualTo(120); + }); + }); + } + @SuppressWarnings("unchecked") @Test void listenerProperties() { @@ -807,6 +826,7 @@ class KafkaAutoConfigurationTests { DefaultKafkaProducerFactory producerFactory = context.getBean(DefaultKafkaProducerFactory.class); DefaultKafkaConsumerFactory consumerFactory = context.getBean(DefaultKafkaConsumerFactory.class); KafkaTemplate kafkaTemplate = context.getBean(KafkaTemplate.class); + assertThat(kafkaTemplate.getKafkaAdmin()).isNull(); assertThat(kafkaTemplate).hasFieldOrPropertyWithValue("producerFactory", producerFactory); AbstractKafkaListenerContainerFactory kafkaListenerContainerFactory = (AbstractKafkaListenerContainerFactory) context .getBean(KafkaListenerContainerFactory.class); @@ -842,6 +862,25 @@ class KafkaAutoConfigurationTests { }); } + @Test + void listenerPropertiesWithAdmin() { + this.contextRunner + .withPropertyValues("spring.kafka.listener.admin.client-id=listener", + "spring.kafka.listener.admin.operation-timeout=2m", + "spring.kafka.listener.admin.properties.fiz.buz=fix.fox") + .run((context) -> { + ConcurrentKafkaListenerContainerFactory factory = context + .getBean(ConcurrentKafkaListenerContainerFactory.class); + ConcurrentMessageListenerContainer container = factory.createContainer("someTopic"); + assertThat(container.getKafkaAdmin()).satisfies((kafkaAdmin) -> { + assertThat(kafkaAdmin.getConfigurationProperties()) + .containsEntry(AdminClientConfig.CLIENT_ID_CONFIG, "listener") + .containsEntry("fiz.buz", "fix.fox"); + assertThat(kafkaAdmin.getOperationTimeout()).isEqualTo(120); + }); + }); + } + @Test void testKafkaTemplateRecordMessageConverters() { this.contextRunner.withUserConfiguration(MessageConverterConfiguration.class) diff --git a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilderTests.java b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilderTests.java index f37ba2ba38e..6663c823978 100644 --- a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilderTests.java +++ b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilderTests.java @@ -35,6 +35,7 @@ import org.junit.jupiter.api.Test; import org.springframework.boot.context.properties.source.MutuallyExclusiveConfigurationPropertiesException; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.IsolationLevel; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Security; +import org.springframework.boot.kafka.autoconfigure.KafkaProperties.SimpleAdmin; import org.springframework.boot.ssl.SslBundle; import org.springframework.core.io.ClassPathResource; import org.springframework.util.unit.DataSize; @@ -264,6 +265,17 @@ class KafkaConfigBuilderTests { .containsEntry("admin.custom", "value"); } + @Test + void adminOverloadUsesGivenAdminProperties() { + KafkaProperties properties = new KafkaProperties(); + SimpleAdmin adminProperties = new SimpleAdmin(); + adminProperties.setClientId("other-admin"); + adminProperties.getProperties().put("other.custom", "value"); + Map config = KafkaConfigBuilder.of(properties).admin(adminProperties).build(); + assertThat(config).containsEntry(AdminClientConfig.CLIENT_ID_CONFIG, "other-admin") + .containsEntry("other.custom", "value"); + } + } @Nested