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
This commit is contained in:
Stéphane Nicoll
2026-09-15 15:39:27 +02:00
parent d031631709
commit dcbb307d3e
9 changed files with 221 additions and 49 deletions
@@ -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);
@@ -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<String, Object> 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;
}
}
@@ -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<Object, Object> recordFilterStrategy;
private final KafkaConnectionDetails kafkaConnectionDetails;
private final BatchMessageConverter batchMessageConverter;
private final @Nullable KafkaTemplate<Object, Object> kafkaTemplate;
@@ -92,7 +96,7 @@ class KafkaAnnotationDrivenConfiguration {
KafkaAnnotationDrivenConfiguration(KafkaProperties properties,
ObjectProvider<RecordMessageConverter> recordMessageConverter,
ObjectProvider<RecordFilterStrategy<Object, Object>> recordFilterStrategy,
ObjectProvider<BatchMessageConverter> batchMessageConverter,
KafkaConnectionDetails kafkaConnectionDetails, ObjectProvider<BatchMessageConverter> batchMessageConverter,
ObjectProvider<KafkaTemplate<Object, Object>> kafkaTemplate,
ObjectProvider<KafkaAwareTransactionManager<Object, Object>> kafkaTransactionManager,
ObjectProvider<ConsumerAwareRebalanceListener> 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(
@@ -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<Object, Object> kafkaProducerFactory,
ProducerListener<Object, Object> kafkaProducerListener,
ObjectProvider<RecordMessageConverter> messageConverter,
ObjectProvider<KafkaTemplateObservationConvention> observationConvention) {
ObjectProvider<KafkaTemplateObservationConvention> observationConvention,
KafkaConnectionDetails kafkaConnectionDetails) {
PropertyMapper map = PropertyMapper.get();
KafkaTemplate<Object, Object> 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<String, Object> 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;
}
@@ -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;
@@ -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<String, String> 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;
}
@@ -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<Object, Object> container = this.factory.createContainer("test");
assertThat(container.getKafkaAdmin()).isSameAs(kafkaAdmin);
}
@Test
@@ -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)
@@ -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<String, Object> config = KafkaConfigBuilder.of(properties).admin(adminProperties).build();
assertThat(config).containsEntry(AdminClientConfig.CLIENT_ID_CONFIG, "other-admin")
.containsEntry("other.custom", "value");
}
}
@Nested