From d031631709a4fca8cb74091c18ce21bc64591a97 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?St=C3=A9phane=20Nicoll?= Date: Thu, 10 Sep 2026 16:04:11 +0200 Subject: [PATCH] Introduce KafkaConfigBuilder This commit introduces a centralized way of building Kafka config properties, and deprecate the equivalent in KafkaProperties. This also removes the customization done in auto-configuration so that KafkaConfigBuilder is the sole source mapping. KafkaConfigBuilder also takes care of applying ConnectionDetails if necessary. Closes gh-51769 --- .../KafkaAnnotationDrivenConfiguration.java | 4 +- .../autoconfigure/KafkaAutoConfiguration.java | 62 +-- .../autoconfigure/KafkaConfigBuilder.java | 481 ++++++++++++++++++ .../kafka/autoconfigure/KafkaProperties.java | 110 ++-- ...aStreamsAnnotationDrivenConfiguration.java | 15 +- .../KafkaConfigBuilderTests.java | 402 +++++++++++++++ .../autoconfigure/KafkaPropertiesTests.java | 86 ---- 7 files changed, 977 insertions(+), 183 deletions(-) create mode 100644 module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilder.java create mode 100644 module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilderTests.java 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 19d52507e22..ffecc542e62 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 @@ -161,8 +161,8 @@ class KafkaAnnotationDrivenConfiguration { ObjectProvider> kafkaConsumerFactory, ObjectProvider>> kafkaContainerCustomizer) { ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); - configurer.configure(factory, kafkaConsumerFactory - .getIfAvailable(() -> new DefaultKafkaConsumerFactory<>(this.properties.buildConsumerProperties()))); + configurer.configure(factory, kafkaConsumerFactory.getIfAvailable( + () -> new DefaultKafkaConsumerFactory<>(KafkaConfigBuilder.of(this.properties).consumer().build()))); kafkaContainerCustomizer.ifAvailable(factory::setContainerCustomizer); return factory; } 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 167ea7758d4..41f359e4f48 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 @@ -21,10 +21,6 @@ import java.time.Duration; import java.util.Map; import java.util.function.Predicate; -import org.apache.kafka.clients.CommonClientConfigs; -import org.apache.kafka.clients.consumer.ConsumerConfig; -import org.apache.kafka.clients.producer.ProducerConfig; -import org.apache.kafka.common.config.SslConfigs; import org.jspecify.annotations.Nullable; import org.springframework.aot.hint.MemberCategory; @@ -40,11 +36,9 @@ 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.KafkaConnectionDetails.Configuration; 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.Template; -import org.springframework.boot.ssl.SslBundle; import org.springframework.boot.ssl.SslBundles; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Import; @@ -64,7 +58,6 @@ import org.springframework.kafka.support.ProducerListener; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.kafka.support.micrometer.KafkaTemplateObservationConvention; import org.springframework.kafka.transaction.KafkaTransactionManager; -import org.springframework.util.StringUtils; import org.springframework.util.backoff.BackOff; /** @@ -131,8 +124,10 @@ public final class KafkaAutoConfiguration { @ConditionalOnMissingBean(ConsumerFactory.class) DefaultKafkaConsumerFactory kafkaConsumerFactory(KafkaConnectionDetails connectionDetails, ObjectProvider customizers) { - Map properties = this.properties.buildConsumerProperties(); - applyKafkaConnectionDetailsForConsumer(properties, connectionDetails); + Map properties = KafkaConfigBuilder.of(this.properties) + .consumer() + .withConnectionDetails(connectionDetails) + .build(); DefaultKafkaConsumerFactory factory = new DefaultKafkaConsumerFactory<>(properties); customizers.orderedStream().forEach((customizer) -> customizer.customize(factory)); return factory; @@ -142,8 +137,10 @@ public final class KafkaAutoConfiguration { @ConditionalOnMissingBean(ProducerFactory.class) DefaultKafkaProducerFactory kafkaProducerFactory(KafkaConnectionDetails connectionDetails, ObjectProvider customizers) { - Map properties = this.properties.buildProducerProperties(); - applyKafkaConnectionDetailsForProducer(properties, connectionDetails); + Map properties = KafkaConfigBuilder.of(this.properties) + .producer() + .withConnectionDetails(connectionDetails) + .build(); DefaultKafkaProducerFactory factory = new DefaultKafkaProducerFactory<>(properties); String transactionIdPrefix = this.properties.getProducer().getTransactionIdPrefix(); if (transactionIdPrefix != null) { @@ -179,8 +176,10 @@ public final class KafkaAutoConfiguration { @Bean @ConditionalOnMissingBean KafkaAdmin kafkaAdmin(KafkaConnectionDetails connectionDetails) { - Map properties = this.properties.buildAdminProperties(); - applyKafkaConnectionDetailsForAdmin(properties, 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) { @@ -209,30 +208,6 @@ public final class KafkaAutoConfiguration { return builder.create(kafkaTemplate); } - private void applyKafkaConnectionDetailsForConsumer(Map properties, - KafkaConnectionDetails connectionDetails) { - Configuration consumer = connectionDetails.getConsumer(); - properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, consumer.getBootstrapServers()); - applySecurityProtocol(properties, consumer.getSecurityProtocol()); - applySslBundle(properties, consumer.getSslBundle()); - } - - private void applyKafkaConnectionDetailsForProducer(Map properties, - KafkaConnectionDetails connectionDetails) { - Configuration producer = connectionDetails.getProducer(); - properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, producer.getBootstrapServers()); - applySecurityProtocol(properties, producer.getSecurityProtocol()); - applySslBundle(properties, producer.getSslBundle()); - } - - private void applyKafkaConnectionDetailsForAdmin(Map properties, - KafkaConnectionDetails connectionDetails) { - Configuration admin = connectionDetails.getAdmin(); - properties.put(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, admin.getBootstrapServers()); - applySecurityProtocol(properties, admin.getSecurityProtocol()); - applySslBundle(properties, admin.getSslBundle()); - } - static BackOff getBackOff(Backoff retryTopicBackoff) { PropertyMapper map = PropertyMapper.get(); RetryPolicy.Builder builder = RetryPolicy.builder().maxRetries(Long.MAX_VALUE); @@ -243,19 +218,6 @@ public final class KafkaAutoConfiguration { return builder.build().getBackOff(); } - static void applySslBundle(Map properties, @Nullable SslBundle sslBundle) { - if (sslBundle != null) { - properties.put(SslConfigs.SSL_ENGINE_FACTORY_CLASS_CONFIG, SslBundleSslEngineFactory.class); - properties.put(SslBundle.class.getName(), sslBundle); - } - } - - static void applySecurityProtocol(Map properties, @Nullable String securityProtocol) { - if (StringUtils.hasLength(securityProtocol)) { - properties.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol); - } - } - static class KafkaRuntimeHints implements RuntimeHintsRegistrar { @Override 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 new file mode 100644 index 00000000000..de1fd4bdf6c --- /dev/null +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilder.java @@ -0,0 +1,481 @@ +/* + * 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.io.IOException; +import java.time.Duration; +import java.util.LinkedHashMap; +import java.util.Locale; +import java.util.Map; +import java.util.function.Consumer; +import java.util.function.Supplier; + +import org.apache.kafka.clients.CommonClientConfigs; +import org.apache.kafka.clients.admin.AdminClientConfig; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.common.config.SslConfigs; +import org.apache.kafka.streams.StreamsConfig; +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; +import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Streams; +import org.springframework.boot.ssl.SslBundle; +import org.springframework.core.io.Resource; +import org.springframework.util.StringUtils; +import org.springframework.util.unit.DataSize; + +/** + * Builder for Kafka components based on {@link KafkaProperties}. + * + * @author Stephane Nicoll + * @since 4.2.0 + */ +public class KafkaConfigBuilder { + + private static final String PROPERTIES_NAMESPACE = "spring.kafka"; + + private final KafkaProperties kafkaProperties; + + protected KafkaConfigBuilder(KafkaProperties kafkaProperties) { + this.kafkaProperties = kafkaProperties; + } + + /** + * Create a builder using the given {@link KafkaProperties}. + * @param kafkaProperties the properties to use + * @return a new builder instance + */ + public static KafkaConfigBuilder of(KafkaProperties kafkaProperties) { + return new KafkaConfigBuilder(kafkaProperties); + } + + /** + * Return a builder for Admin-related configuration. + * @return an admin config builder + * @see AdminClientConfig + */ + public ConfigBuilder admin() { + return new AdminConfigBuilder(initializeKafkaConfig(), this.kafkaProperties.getAdmin(), null); + } + + /** + * Return a builder for Consumer-related configuration. + * @return a consumer config builder + * @see ConsumerConfig + */ + public ConfigBuilder consumer() { + return new ConsumerConfigBuilder(initializeKafkaConfig(), this.kafkaProperties.getConsumer(), null); + } + + /** + * Return a builder for Producer-related configuration. + * @return a producer config builder + * @see ProducerConfig + */ + public ConfigBuilder producer() { + return new ProducerConfigBuilder(initializeKafkaConfig(), this.kafkaProperties.getProducer(), null); + } + + /** + * Return a builder for Streams-related configuration. + * @return a streams config builder + * @see StreamsConfig + */ + public ConfigBuilder streams() { + return new StreamsConfigBuilder(initializeKafkaConfig(), this.kafkaProperties.getStreams(), null); + } + + protected KafkaConfig initializeKafkaConfig() { + KafkaConfig kafkaConfig = new KafkaConfig(); + kafkaConfig.putIfNonNull(this.kafkaProperties::getBootstrapServers, + CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG); + kafkaConfig.putIfNonNull(this.kafkaProperties::getClientId, CommonClientConfigs.CLIENT_ID_CONFIG); + new SslConfigBuilder(this.kafkaProperties.getSsl(), PROPERTIES_NAMESPACE + ".ssl").apply(kafkaConfig); + new SecurityConfigBuilder(this.kafkaProperties.getSecurity()).apply(kafkaConfig); + kafkaConfig.putAll(this.kafkaProperties.getProperties()); + return kafkaConfig; + } + + private static void applySecurityProtocol(Map properties, @Nullable String securityProtocol) { + if (StringUtils.hasLength(securityProtocol)) { + properties.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol); + } + } + + private static void applySslBundle(Map properties, @Nullable SslBundle sslBundle) { + if (sslBundle != null) { + properties.put(SslConfigs.SSL_ENGINE_FACTORY_CLASS_CONFIG, SslBundleSslEngineFactory.class); + properties.put(SslBundle.class.getName(), sslBundle); + } + } + + public interface ConfigBuilder { + + /** + * Apply the given {@link KafkaConnectionDetails}. + * @param connectionDetails the connection details to use + * @return {@code this} + */ + ConfigBuilder withConnectionDetails(@Nullable KafkaConnectionDetails connectionDetails); + + /** + * Build the configuration managed by this builder. + * @return the configuration to use + */ + Map build(); + + } + + protected static class KafkaConfig extends LinkedHashMap { + + public KafkaConfig() { + } + + public KafkaConfig(Map m) { + super(m); + } + + public void putIfNonNull(Supplier<@Nullable Object> value, String key) { + Object toSet = value.get(); + if (toSet != null) { + this.put(key, toSet); + } + } + + public Consumer in(String key) { + return (value) -> put(key, value); + } + + } + + private static final class AdminConfigBuilder implements ConfigBuilder { + + private static final String NAMESPACE = PROPERTIES_NAMESPACE + ".admin"; + + private final KafkaConfig kafkaConfig; + + private final KafkaProperties.Admin admin; + + private final @Nullable KafkaConnectionDetails connectionDetails; + + private AdminConfigBuilder(KafkaConfig kafkaConfig, Admin admin, + @Nullable KafkaConnectionDetails connectionDetails) { + this.kafkaConfig = kafkaConfig; + this.admin = admin; + this.connectionDetails = connectionDetails; + } + + @Override + public ConfigBuilder withConnectionDetails(@Nullable KafkaConnectionDetails connectionDetails) { + return new AdminConfigBuilder(this.kafkaConfig, this.admin, connectionDetails); + } + + @Override + public Map build() { + KafkaConfig adminConfig = new KafkaConfig(this.kafkaConfig); + PropertyMapper map = PropertyMapper.get(); + map.from(this.admin::getClientId).to(adminConfig.in(ProducerConfig.CLIENT_ID_CONFIG)); + new SslConfigBuilder(this.admin.getSsl(), NAMESPACE + ".ssl").apply(adminConfig); + new SecurityConfigBuilder(this.admin.getSecurity()).apply(adminConfig); + adminConfig.putAll(this.admin.getProperties()); + if (this.connectionDetails != null) { + Configuration admin = this.connectionDetails.getAdmin(); + adminConfig.put(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, admin.getBootstrapServers()); + applySecurityProtocol(adminConfig, admin.getSecurityProtocol()); + applySslBundle(adminConfig, admin.getSslBundle()); + } + return adminConfig; + } + + } + + private static final class ConsumerConfigBuilder implements ConfigBuilder { + + private static final String NAMESPACE = PROPERTIES_NAMESPACE + ".consumer"; + + private final KafkaConfig kafkaConfig; + + private final KafkaProperties.Consumer consumer; + + private final @Nullable KafkaConnectionDetails connectionDetails; + + private ConsumerConfigBuilder(KafkaConfig kafkaConfig, KafkaProperties.Consumer consumer, + @Nullable KafkaConnectionDetails connectionDetails) { + this.kafkaConfig = kafkaConfig; + this.consumer = consumer; + this.connectionDetails = connectionDetails; + } + + @Override + public ConsumerConfigBuilder withConnectionDetails(@Nullable KafkaConnectionDetails connectionDetails) { + return new ConsumerConfigBuilder(this.kafkaConfig, this.consumer, connectionDetails); + } + + @Override + public Map build() { + KafkaConfig consumerConfig = new KafkaConfig(this.kafkaConfig); + PropertyMapper map = PropertyMapper.get(); + map.from(this.consumer::getAutoCommitInterval) + .asInt(Duration::toMillis) + .to(consumerConfig.in(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG)); + map.from(this.consumer::getAutoOffsetReset).to(consumerConfig.in(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)); + map.from(this.consumer::getBootstrapServers).to(consumerConfig.in(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG)); + map.from(this.consumer::getClientId).to(consumerConfig.in(ConsumerConfig.CLIENT_ID_CONFIG)); + map.from(this.consumer::getEnableAutoCommit) + .to(consumerConfig.in(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG)); + map.from(this.consumer::getFetchMaxWait) + .asInt(Duration::toMillis) + .to(consumerConfig.in(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG)); + map.from(this.consumer::getFetchMinSize) + .asInt(DataSize::toBytes) + .to(consumerConfig.in(ConsumerConfig.FETCH_MIN_BYTES_CONFIG)); + map.from(this.consumer::getGroupId).to(consumerConfig.in(ConsumerConfig.GROUP_ID_CONFIG)); + map.from(this.consumer::getHeartbeatInterval) + .asInt(Duration::toMillis) + .to(consumerConfig.in(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG)); + map.from(() -> this.consumer.getIsolationLevel().name().toLowerCase(Locale.ROOT)) + .to(consumerConfig.in(ConsumerConfig.ISOLATION_LEVEL_CONFIG)); + map.from(this.consumer::getKeyDeserializer) + .to(consumerConfig.in(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG)); + map.from(this.consumer::getValueDeserializer) + .to(consumerConfig.in(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG)); + map.from(this.consumer::getMaxPollRecords).to(consumerConfig.in(ConsumerConfig.MAX_POLL_RECORDS_CONFIG)); + map.from(this.consumer::getMaxPollInterval) + .asInt(Duration::toMillis) + .to(consumerConfig.in(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG)); + new SslConfigBuilder(this.consumer.getSsl(), NAMESPACE + ".ssl").apply(consumerConfig); + new SecurityConfigBuilder(this.consumer.getSecurity()).apply(consumerConfig); + consumerConfig.putAll(this.consumer.getProperties()); + if (this.connectionDetails != null) { + Configuration consumer = this.connectionDetails.getConsumer(); + consumerConfig.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, consumer.getBootstrapServers()); + applySecurityProtocol(consumerConfig, consumer.getSecurityProtocol()); + applySslBundle(consumerConfig, consumer.getSslBundle()); + } + return consumerConfig; + } + + } + + private static final class ProducerConfigBuilder implements ConfigBuilder { + + private static final String NAMESPACE = PROPERTIES_NAMESPACE + ".producer"; + + private final KafkaConfig kafkaConfig; + + private final KafkaProperties.Producer producer; + + private final @Nullable KafkaConnectionDetails connectionDetails; + + private ProducerConfigBuilder(KafkaConfig kafkaConfig, Producer producer, + @Nullable KafkaConnectionDetails connectionDetails) { + this.producer = producer; + this.kafkaConfig = kafkaConfig; + this.connectionDetails = connectionDetails; + } + + @Override + public ProducerConfigBuilder withConnectionDetails(@Nullable KafkaConnectionDetails connectionDetails) { + return new ProducerConfigBuilder(this.kafkaConfig, this.producer, connectionDetails); + } + + @Override + public Map build() { + KafkaConfig producerConfig = new KafkaConfig(this.kafkaConfig); + PropertyMapper map = PropertyMapper.get(); + map.from(this.producer::getAcks).to(producerConfig.in(ProducerConfig.ACKS_CONFIG)); + map.from(this.producer::getBatchSize) + .asInt(DataSize::toBytes) + .to(producerConfig.in(ProducerConfig.BATCH_SIZE_CONFIG)); + map.from(this.producer::getBootstrapServers).to(producerConfig.in(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG)); + map.from(this.producer::getBufferMemory) + .as(DataSize::toBytes) + .to(producerConfig.in(ProducerConfig.BUFFER_MEMORY_CONFIG)); + map.from(this.producer::getClientId).to(producerConfig.in(ProducerConfig.CLIENT_ID_CONFIG)); + map.from(this.producer::getCompressionType).to(producerConfig.in(ProducerConfig.COMPRESSION_TYPE_CONFIG)); + map.from(this.producer::getKeySerializer).to(producerConfig.in(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG)); + map.from(this.producer::getRetries).to(producerConfig.in(ProducerConfig.RETRIES_CONFIG)); + map.from(this.producer::getValueSerializer) + .to(producerConfig.in(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG)); + new SslConfigBuilder(this.producer.getSsl(), NAMESPACE + ".ssl").apply(producerConfig); + new SecurityConfigBuilder(this.producer.getSecurity()).apply(producerConfig); + producerConfig.putAll(this.producer.getProperties()); + if (this.connectionDetails != null) { + Configuration producer = this.connectionDetails.getProducer(); + producerConfig.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, producer.getBootstrapServers()); + applySecurityProtocol(producerConfig, producer.getSecurityProtocol()); + applySslBundle(producerConfig, producer.getSslBundle()); + } + return producerConfig; + } + + } + + private static final class StreamsConfigBuilder implements ConfigBuilder { + + private static final String NAMESPACE = PROPERTIES_NAMESPACE + ".streams"; + + private final KafkaConfig kafkaConfig; + + private final Streams streams; + + private final @Nullable KafkaConnectionDetails connectionDetails; + + private StreamsConfigBuilder(KafkaConfig kafkaConfig, Streams streams, + @Nullable KafkaConnectionDetails connectionDetails) { + this.kafkaConfig = kafkaConfig; + this.streams = streams; + this.connectionDetails = connectionDetails; + } + + @Override + public StreamsConfigBuilder withConnectionDetails(@Nullable KafkaConnectionDetails connectionDetails) { + return new StreamsConfigBuilder(this.kafkaConfig, this.streams, connectionDetails); + } + + @Override + public Map build() { + KafkaConfig streamsConfig = new KafkaConfig(this.kafkaConfig); + PropertyMapper map = PropertyMapper.get(); + map.from(this.streams::getApplicationId).to(streamsConfig.in("application.id")); + map.from(this.streams::getBootstrapServers).to(streamsConfig.in(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG)); + map.from(this.streams::getStateStoreCacheMaxSize) + .asInt(DataSize::toBytes) + .to(streamsConfig.in("statestore.cache.max.bytes")); + map.from(this.streams::getClientId).to(streamsConfig.in(StreamsConfig.CLIENT_ID_CONFIG)); + map.from(this.streams::getReplicationFactor).to(streamsConfig.in("replication.factor")); + map.from(this.streams::getStateDir).to(streamsConfig.in("state.dir")); + new SslConfigBuilder(this.streams.getSsl(), NAMESPACE + ".ssl").apply(streamsConfig); + new SecurityConfigBuilder(this.streams.getSecurity()).apply(streamsConfig); + streamsConfig.putAll(this.streams.getProperties()); + if (this.connectionDetails != null) { + Configuration streams = this.connectionDetails.getStreams(); + streamsConfig.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, streams.getBootstrapServers()); + applySecurityProtocol(streamsConfig, streams.getSecurityProtocol()); + applySslBundle(streamsConfig, streams.getSslBundle()); + } + return streamsConfig; + } + + } + + private static final class SslConfigBuilder { + + private final Ssl ssl; + + private final String prefix; + + private SslConfigBuilder(Ssl ssl, String prefix) { + this.ssl = ssl; + this.prefix = prefix; + } + + void apply(KafkaConfig kafkaConfig) { + validate(); + String bundleName = this.ssl.getBundle(); + if (StringUtils.hasText(bundleName)) { + return; + } + PropertyMapper map = PropertyMapper.get(); + map.from(this.ssl::getKeyPassword).to(kafkaConfig.in(SslConfigs.SSL_KEY_PASSWORD_CONFIG)); + map.from(this.ssl::getKeyStoreCertificateChain) + .to(kafkaConfig.in(SslConfigs.SSL_KEYSTORE_CERTIFICATE_CHAIN_CONFIG)); + map.from(this.ssl::getKeyStoreKey).to(kafkaConfig.in(SslConfigs.SSL_KEYSTORE_KEY_CONFIG)); + map.from(this.ssl::getKeyStoreLocation) + .as(this::resourceToPath) + .to(kafkaConfig.in(SslConfigs.SSL_KEYSTORE_LOCATION_CONFIG)); + map.from(this.ssl::getKeyStorePassword).to(kafkaConfig.in(SslConfigs.SSL_KEYSTORE_PASSWORD_CONFIG)); + map.from(this.ssl::getKeyStoreType).to(kafkaConfig.in(SslConfigs.SSL_KEYSTORE_TYPE_CONFIG)); + map.from(this.ssl::getTrustStoreCertificates) + .to(kafkaConfig.in(SslConfigs.SSL_TRUSTSTORE_CERTIFICATES_CONFIG)); + map.from(this.ssl::getTrustStoreLocation) + .as(this::resourceToPath) + .to(kafkaConfig.in(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG)); + map.from(this.ssl::getTrustStorePassword).to(kafkaConfig.in(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG)); + map.from(this.ssl::getTrustStoreType).to(kafkaConfig.in(SslConfigs.SSL_TRUSTSTORE_TYPE_CONFIG)); + map.from(this.ssl::getProtocol).to(kafkaConfig.in(SslConfigs.SSL_PROTOCOL_CONFIG)); + } + + private void validate() { + MutuallyExclusiveConfigurationPropertiesException.throwIfMultipleMatchingValuesIn((entries) -> { + entries.put(key("key-store-key"), this.ssl.getKeyStoreKey()); + entries.put(key("key-store-location"), this.ssl.getKeyStoreLocation()); + }, this::hasValue); + MutuallyExclusiveConfigurationPropertiesException.throwIfMultipleMatchingValuesIn((entries) -> { + entries.put(key("trust-store-certificates"), this.ssl.getTrustStoreCertificates()); + entries.put(key("trust-store-location"), this.ssl.getTrustStoreLocation()); + }, this::hasValue); + MutuallyExclusiveConfigurationPropertiesException.throwIfMultipleMatchingValuesIn((entries) -> { + entries.put(key("bundle"), this.ssl.getBundle()); + entries.put(key("key-store-key"), this.ssl.getKeyStoreKey()); + }, this::hasValue); + MutuallyExclusiveConfigurationPropertiesException.throwIfMultipleMatchingValuesIn((entries) -> { + entries.put(key("bundle"), this.ssl.getBundle()); + entries.put(key("key-store-location"), this.ssl.getKeyStoreLocation()); + }, this::hasValue); + MutuallyExclusiveConfigurationPropertiesException.throwIfMultipleMatchingValuesIn((entries) -> { + entries.put(key("bundle"), this.ssl.getBundle()); + entries.put(key("trust-store-certificates"), this.ssl.getTrustStoreCertificates()); + }, this::hasValue); + MutuallyExclusiveConfigurationPropertiesException.throwIfMultipleMatchingValuesIn((entries) -> { + entries.put(key("bundle"), this.ssl.getBundle()); + entries.put(key("trust-store-location"), this.ssl.getTrustStoreLocation()); + }, this::hasValue); + } + + private String key(String name) { + return "%s.%s".formatted(this.prefix, name); + } + + private boolean hasValue(@Nullable Object value) { + return (value instanceof String string) ? StringUtils.hasText(string) : value != null; + } + + private String resourceToPath(Resource resource) { + try { + return resource.getFile().getAbsolutePath(); + } + catch (IOException ex) { + throw new IllegalStateException("Resource '" + resource + "' must be on a file system", ex); + } + } + + } + + private static final class SecurityConfigBuilder { + + private final Security security; + + SecurityConfigBuilder(Security security) { + this.security = security; + } + + void apply(KafkaConfig kafkaConfig) { + PropertyMapper map = PropertyMapper.get(); + map.from(this.security::getProtocol).to(kafkaConfig.in(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG)); + } + + } + +} 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 be2af83fb7b..8ae3d38042b 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 @@ -41,7 +41,6 @@ import org.springframework.boot.convert.DurationUnit; import org.springframework.core.io.Resource; import org.springframework.kafka.listener.ContainerProperties.AckMode; import org.springframework.kafka.security.jaas.KafkaJaasLoginModuleInitializer; -import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; import org.springframework.util.unit.DataSize; @@ -161,20 +160,15 @@ public class KafkaProperties { return this.retry; } + /** + * Create an initial map of properties for common settings. + * @return the consumer properties initialized with the customizations defined on this + * instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of {@link KafkaConfigBuilder} + */ + @Deprecated(since = "4.2.0", forRemoval = true) private Map buildCommonProperties() { - Map properties = new LinkedHashMap<>(); - if (this.bootstrapServers != null) { - properties.put(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, this.bootstrapServers); - } - if (this.clientId != null) { - properties.put(CommonClientConfigs.CLIENT_ID_CONFIG, this.clientId); - } - properties.putAll(this.ssl.buildProperties()); - properties.putAll(this.security.buildProperties()); - if (!CollectionUtils.isEmpty(this.properties)) { - properties.putAll(this.properties); - } - return properties; + return KafkaConfigBuilder.of(this).initializeKafkaConfig(); } /** @@ -182,13 +176,13 @@ public class KafkaProperties { *

* This allows you to add additional properties, if necessary, and override the * default {@code kafkaConsumerFactory} bean. - * @return the consumer properties initialized with the customizations defined on this - * instance + * @return common and consumer properties initialized with the customizations defined + * on this instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of {@link KafkaConfigBuilder} */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildConsumerProperties() { - Map properties = buildCommonProperties(); - properties.putAll(this.consumer.buildProperties()); - return properties; + return KafkaConfigBuilder.of(this).consumer().build(); } /** @@ -196,13 +190,13 @@ public class KafkaProperties { *

* This allows you to add additional properties, if necessary, and override the * default {@code kafkaProducerFactory} bean. - * @return the producer properties initialized with the customizations defined on this - * instance + * @return common and producer properties initialized with the customizations defined + * on this instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of {@link KafkaConfigBuilder} */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildProducerProperties() { - Map properties = buildCommonProperties(); - properties.putAll(this.producer.buildProperties()); - return properties; + return KafkaConfigBuilder.of(this).producer().build(); } /** @@ -210,26 +204,26 @@ public class KafkaProperties { *

* This allows you to add additional properties, if necessary, and override the * default {@code kafkaAdmin} bean. - * @return the admin properties initialized with the customizations defined on this - * instance + * @return common and admin properties initialized with the customizations defined on + * this instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of {@link KafkaConfigBuilder} */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildAdminProperties() { - Map properties = buildCommonProperties(); - properties.putAll(this.admin.buildProperties()); - return properties; + return KafkaConfigBuilder.of(this).admin().build(); } /** * Create an initial map of streams properties from the state of this instance. *

* This allows you to add additional properties, if necessary. - * @return the streams properties initialized with the customizations defined on this - * instance + * @return common and streams properties initialized with the customizations defined + * on this instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of {@link KafkaConfigBuilder} */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildStreamsProperties() { - Map properties = buildCommonProperties(); - properties.putAll(this.streams.buildProperties()); - return properties; + return KafkaConfigBuilder.of(this).streams().build(); } public static class Consumer { @@ -444,6 +438,14 @@ public class KafkaProperties { return this.properties; } + /** + * Build the relevant consumer properties. + * @return the consumer properties initialized with the customizations defined on + * this instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of + * {@link KafkaConfigBuilder} + */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildProperties() { Properties properties = new Properties(); PropertyMapper map = PropertyMapper.get(); @@ -634,6 +636,14 @@ public class KafkaProperties { return this.properties; } + /** + * Build the relevant producer properties. + * @return the producer properties initialized with the customizations defined on + * this instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of + * {@link KafkaConfigBuilder} + */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildProperties() { Properties properties = new Properties(); PropertyMapper map = PropertyMapper.get(); @@ -755,6 +765,14 @@ public class KafkaProperties { return this.properties; } + /** + * Build the relevant admin properties. + * @return the admin properties initialized with the customizations defined on + * this instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of + * {@link KafkaConfigBuilder} + */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildProperties() { Properties properties = new Properties(); PropertyMapper map = PropertyMapper.get(); @@ -902,6 +920,14 @@ public class KafkaProperties { return this.properties; } + /** + * Build the relevant stream properties. + * @return the stream properties initialized with the customizations defined on + * this instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of + * {@link KafkaConfigBuilder} + */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildProperties() { Properties properties = new Properties(); PropertyMapper map = PropertyMapper.get(); @@ -1433,6 +1459,14 @@ public class KafkaProperties { this.protocol = protocol; } + /** + * Build the relevant SSL properties. + * @return the SSL properties initialized with the customizations defined on this + * instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of + * {@link KafkaConfigBuilder} + */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildProperties() { validate(); String bundleName = getBundle(); @@ -1575,6 +1609,14 @@ public class KafkaProperties { this.protocol = protocol; } + /** + * Build the relevant security properties. + * @return the security properties initialized with the customizations defined on + * this instance + * @deprecated since 4.2.0 for removal in 4.4.0 in favor of + * {@link KafkaConfigBuilder} + */ + @Deprecated(since = "4.2.0", forRemoval = true) public Map buildProperties() { Properties properties = new Properties(); PropertyMapper map = PropertyMapper.get(); 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 410d7206725..c5d5ab301e7 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 @@ -18,7 +18,6 @@ package org.springframework.boot.kafka.autoconfigure; import java.util.Map; -import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; @@ -61,8 +60,10 @@ class KafkaStreamsAnnotationDrivenConfiguration { @Bean(KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME) KafkaStreamsConfiguration defaultKafkaStreamsConfig(Environment environment, KafkaConnectionDetails connectionDetails) { - Map properties = this.properties.buildStreamsProperties(); - applyKafkaConnectionDetailsForStreams(properties, connectionDetails); + Map properties = KafkaConfigBuilder.of(this.properties) + .streams() + .withConnectionDetails(connectionDetails) + .build(); if (this.properties.getStreams().getApplicationId() == null) { String applicationName = environment.getProperty("spring.application.name"); if (applicationName == null) { @@ -79,14 +80,6 @@ class KafkaStreamsAnnotationDrivenConfiguration { return new KafkaPropertiesStreamsBuilderFactoryBeanConfigurer(this.properties); } - private void applyKafkaConnectionDetailsForStreams(Map properties, - KafkaConnectionDetails connectionDetails) { - KafkaConnectionDetails.Configuration streams = connectionDetails.getStreams(); - properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, streams.getBootstrapServers()); - KafkaAutoConfiguration.applySecurityProtocol(properties, streams.getSecurityProtocol()); - KafkaAutoConfiguration.applySslBundle(properties, streams.getSslBundle()); - } - static class KafkaPropertiesStreamsBuilderFactoryBeanConfigurer implements StreamsBuilderFactoryBeanConfigurer { private final KafkaProperties properties; 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 new file mode 100644 index 00000000000..f37ba2ba38e --- /dev/null +++ b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaConfigBuilderTests.java @@ -0,0 +1,402 @@ +/* + * 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.time.Duration; +import java.util.List; +import java.util.Map; + +import org.apache.kafka.clients.CommonClientConfigs; +import org.apache.kafka.clients.admin.AdminClientConfig; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.common.config.SslConfigs; +import org.apache.kafka.common.serialization.LongDeserializer; +import org.apache.kafka.common.serialization.LongSerializer; +import org.apache.kafka.streams.StreamsConfig; +import org.jspecify.annotations.Nullable; +import org.junit.jupiter.api.Nested; +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.ssl.SslBundle; +import org.springframework.core.io.ClassPathResource; +import org.springframework.util.unit.DataSize; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatExceptionOfType; +import static org.mockito.Mockito.mock; + +/** + * Tests for {@link KafkaConfigBuilder}. + * + * @author Stephane Nicoll + */ +class KafkaConfigBuilderTests { + + private static KafkaConnectionDetails connectionDetails(List bootstrapServers, + @Nullable String securityProtocol, @Nullable SslBundle sslBundle) { + return new KafkaConnectionDetails() { + + @Override + public List getBootstrapServers() { + return bootstrapServers; + } + + @Override + public @Nullable String getSecurityProtocol() { + return securityProtocol; + } + + @Override + public @Nullable SslBundle getSslBundle() { + return sslBundle; + } + + }; + } + + interface CommonTests { + + Map build(KafkaProperties properties, @Nullable KafkaConnectionDetails connectionDetails); + + Security security(KafkaProperties properties); + + default Map build(KafkaProperties properties) { + return build(properties, null); + } + + @Test + default void baseBootstrapServersIsApplied() { + KafkaProperties properties = new KafkaProperties(); + properties.setBootstrapServers(List.of("base:9092")); + Map config = build(properties); + assertThat(config).containsEntry(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, List.of("base:9092")); + } + + @Test + default void baseClientIdIsApplied() { + KafkaProperties properties = new KafkaProperties(); + properties.setClientId("base-client"); + Map config = build(properties); + assertThat(config).containsEntry(CommonClientConfigs.CLIENT_ID_CONFIG, "base-client"); + } + + @Test + default void baseSecurityProtocolIsApplied() { + KafkaProperties properties = new KafkaProperties(); + properties.getSecurity().setProtocol("SASL_PLAINTEXT"); + Map config = build(properties); + assertThat(config).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_PLAINTEXT"); + } + + @Test + default void componentSecurityProtocolOverridesBase() { + KafkaProperties properties = new KafkaProperties(); + properties.getSecurity().setProtocol("PLAINTEXT"); + security(properties).setProtocol("SASL_SSL"); + Map config = build(properties); + assertThat(config).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL"); + } + + @Test + default void baseAdditionalPropertiesAreApplied() { + KafkaProperties properties = new KafkaProperties(); + properties.getProperties().put("base.custom", "value"); + Map config = build(properties); + assertThat(config).containsEntry("base.custom", "value"); + } + + @Test + default void connectionDetailsOverridesBootstrapServers() { + KafkaProperties properties = new KafkaProperties(); + KafkaConnectionDetails connectionDetails = connectionDetails(List.of("details:9092"), null, null); + Map config = build(properties, connectionDetails); + assertThat(config).containsEntry(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, List.of("details:9092")); + } + + @Test + default void connectionDetailsOverridesSecurityProtocol() { + KafkaProperties properties = new KafkaProperties(); + KafkaConnectionDetails connectionDetails = connectionDetails(List.of("localhost:9092"), "SSL", null); + Map config = build(properties, connectionDetails); + assertThat(config).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL"); + } + + @Test + default void connectionDetailsWithoutSecurityProtocolDoesNotOverrideExisting() { + KafkaProperties properties = new KafkaProperties(); + properties.getSecurity().setProtocol("SASL_SSL"); + KafkaConnectionDetails connectionDetails = connectionDetails(List.of("localhost:9092"), null, null); + Map config = build(properties, connectionDetails); + assertThat(config).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL"); + } + + @Test + default void connectionDetailsAppliesSslBundle() { + KafkaProperties properties = new KafkaProperties(); + SslBundle sslBundle = mock(SslBundle.class); + KafkaConnectionDetails connectionDetails = connectionDetails(List.of("localhost:9092"), null, sslBundle); + Map config = build(properties, connectionDetails); + assertThat(config).containsEntry(SslConfigs.SSL_ENGINE_FACTORY_CLASS_CONFIG, + SslBundleSslEngineFactory.class); + assertThat(config).containsEntry(SslBundle.class.getName(), sslBundle); + } + + @Test + default void sslPemConfiguration() { + KafkaProperties properties = new KafkaProperties(); + properties.getSsl().setKeyStoreKey("-----BEGINkey"); + properties.getSsl().setTrustStoreCertificates("-----BEGINtrust"); + properties.getSsl().setKeyStoreCertificateChain("-----BEGINchain"); + Map config = build(properties); + assertThat(config).containsEntry(SslConfigs.SSL_KEYSTORE_KEY_CONFIG, "-----BEGINkey"); + assertThat(config).containsEntry(SslConfigs.SSL_TRUSTSTORE_CERTIFICATES_CONFIG, "-----BEGINtrust"); + assertThat(config).containsEntry(SslConfigs.SSL_KEYSTORE_CERTIFICATE_CHAIN_CONFIG, "-----BEGINchain"); + } + + @Test + default void sslPemConfigurationWithEmptyBundle() { + KafkaProperties properties = new KafkaProperties(); + properties.getSsl().setKeyStoreKey("-----BEGINkey"); + properties.getSsl().setTrustStoreCertificates("-----BEGINtrust"); + properties.getSsl().setKeyStoreCertificateChain("-----BEGINchain"); + properties.getSsl().setBundle(""); + Map config = build(properties); + assertThat(config).containsEntry(SslConfigs.SSL_KEYSTORE_KEY_CONFIG, "-----BEGINkey"); + assertThat(config).containsEntry(SslConfigs.SSL_TRUSTSTORE_CERTIFICATES_CONFIG, "-----BEGINtrust"); + assertThat(config).containsEntry(SslConfigs.SSL_KEYSTORE_CERTIFICATE_CHAIN_CONFIG, "-----BEGINchain"); + } + + @Test + default void sslPropertiesWhenKeyStoreLocationAndKeySetShouldThrowException() { + KafkaProperties properties = new KafkaProperties(); + properties.getSsl().setKeyStoreKey("-----BEGIN"); + properties.getSsl().setKeyStoreLocation(new ClassPathResource("ksLoc")); + assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) + .isThrownBy(() -> build(properties)); + } + + @Test + default void sslPropertiesWhenTrustStoreLocationAndCertificatesSetShouldThrowException() { + KafkaProperties properties = new KafkaProperties(); + properties.getSsl().setTrustStoreLocation(new ClassPathResource("tsLoc")); + properties.getSsl().setTrustStoreCertificates("-----BEGIN"); + assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) + .isThrownBy(() -> build(properties)); + } + + @Test + default void sslPropertiesWhenKeyStoreLocationAndBundleSetShouldThrowException() { + KafkaProperties properties = new KafkaProperties(); + properties.getSsl().setBundle("myBundle"); + properties.getSsl().setKeyStoreLocation(new ClassPathResource("ksLoc")); + assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) + .isThrownBy(() -> build(properties)); + } + + @Test + default void sslPropertiesWhenKeyStoreKeyAndBundleSetShouldThrowException() { + KafkaProperties properties = new KafkaProperties(); + properties.getSsl().setBundle("myBundle"); + properties.getSsl().setKeyStoreKey("-----BEGIN"); + assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) + .isThrownBy(() -> build(properties)); + } + + @Test + default void sslPropertiesWhenTrustStoreLocationAndBundleSetShouldThrowException() { + KafkaProperties properties = new KafkaProperties(); + properties.getSsl().setBundle("myBundle"); + properties.getSsl().setTrustStoreLocation(new ClassPathResource("tsLoc")); + assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) + .isThrownBy(() -> build(properties)); + } + + @Test + default void sslPropertiesWhenTrustStoreCertificatesAndBundleSetShouldThrowException() { + KafkaProperties properties = new KafkaProperties(); + properties.getSsl().setBundle("myBundle"); + properties.getSsl().setTrustStoreCertificates("-----BEGIN"); + assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) + .isThrownBy(() -> build(properties)); + } + + } + + @Nested + class AdminConfigTests implements CommonTests { + + @Override + public Map build(KafkaProperties properties, + @Nullable KafkaConnectionDetails connectionDetails) { + return KafkaConfigBuilder.of(properties).admin().withConnectionDetails(connectionDetails).build(); + } + + @Override + public Security security(KafkaProperties properties) { + return properties.getAdmin().getSecurity(); + } + + @Test + void adminPropertiesAreApplied() { + KafkaProperties properties = new KafkaProperties(); + properties.getAdmin().setClientId("admin-client"); + properties.getAdmin().getProperties().put("admin.custom", "value"); + Map config = build(properties); + assertThat(config).containsEntry(AdminClientConfig.CLIENT_ID_CONFIG, "admin-client") + .containsEntry("admin.custom", "value"); + } + + } + + @Nested + class ConsumerConfigTests implements CommonTests { + + @Override + public Map build(KafkaProperties properties, + @Nullable KafkaConnectionDetails connectionDetails) { + return KafkaConfigBuilder.of(properties).consumer().withConnectionDetails(connectionDetails).build(); + } + + @Override + public Security security(KafkaProperties properties) { + return properties.getConsumer().getSecurity(); + } + + @Test + void consumerPropertiesAreApplied() { + KafkaProperties properties = new KafkaProperties(); + KafkaProperties.Consumer consumer = properties.getConsumer(); + consumer.setAutoCommitInterval(Duration.ofSeconds(5)); + consumer.setAutoOffsetReset("earliest"); + consumer.setBootstrapServers(List.of("consumer:9092")); + consumer.setClientId("consumer-client"); + consumer.setEnableAutoCommit(true); + consumer.setFetchMaxWait(Duration.ofSeconds(2)); + consumer.setFetchMinSize(DataSize.ofKilobytes(2)); + consumer.setGroupId("group"); + consumer.setHeartbeatInterval(Duration.ofSeconds(3)); + consumer.setIsolationLevel(IsolationLevel.READ_COMMITTED); + consumer.setKeyDeserializer(LongDeserializer.class); + consumer.setValueDeserializer(LongDeserializer.class); + consumer.setMaxPollRecords(10); + consumer.setMaxPollInterval(Duration.ofSeconds(30)); + consumer.getProperties().put("consumer.custom", "value"); + Map config = build(properties); + assertThat(config).containsEntry(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 5000) + .containsEntry(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest") + .containsEntry(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, List.of("consumer:9092")) + .containsEntry(ConsumerConfig.CLIENT_ID_CONFIG, "consumer-client") + .containsEntry(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true) + .containsEntry(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG, 2000) + .containsEntry(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, (int) DataSize.ofKilobytes(2).toBytes()) + .containsEntry(ConsumerConfig.GROUP_ID_CONFIG, "group") + .containsEntry(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 3000) + .containsEntry(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed") + .containsEntry(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class) + .containsEntry(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class) + .containsEntry(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 10) + .containsEntry(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 30000) + .containsEntry("consumer.custom", "value"); + } + + } + + @Nested + class ProducerConfigTests implements CommonTests { + + @Override + public Map build(KafkaProperties properties, + @Nullable KafkaConnectionDetails connectionDetails) { + return KafkaConfigBuilder.of(properties).producer().withConnectionDetails(connectionDetails).build(); + } + + @Override + public Security security(KafkaProperties properties) { + return properties.getProducer().getSecurity(); + } + + @Test + void producerPropertiesAreApplied() { + KafkaProperties properties = new KafkaProperties(); + KafkaProperties.Producer producer = properties.getProducer(); + producer.setAcks("all"); + producer.setBatchSize(DataSize.ofKilobytes(16)); + producer.setBootstrapServers(List.of("producer:9092")); + producer.setBufferMemory(DataSize.ofMegabytes(32)); + producer.setClientId("producer-client"); + producer.setCompressionType("gzip"); + producer.setKeySerializer(LongSerializer.class); + producer.setRetries(3); + producer.setValueSerializer(LongSerializer.class); + producer.getProperties().put("producer.custom", "value"); + Map config = build(properties); + assertThat(config).containsEntry(ProducerConfig.ACKS_CONFIG, "all") + .containsEntry(ProducerConfig.BATCH_SIZE_CONFIG, (int) DataSize.ofKilobytes(16).toBytes()) + .containsEntry(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, List.of("producer:9092")) + .containsEntry(ProducerConfig.BUFFER_MEMORY_CONFIG, DataSize.ofMegabytes(32).toBytes()) + .containsEntry(ProducerConfig.CLIENT_ID_CONFIG, "producer-client") + .containsEntry(ProducerConfig.COMPRESSION_TYPE_CONFIG, "gzip") + .containsEntry(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class) + .containsEntry(ProducerConfig.RETRIES_CONFIG, 3) + .containsEntry(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, LongSerializer.class) + .containsEntry("producer.custom", "value"); + } + + } + + @Nested + class StreamsConfigTests implements CommonTests { + + @Override + public Map build(KafkaProperties properties, + @Nullable KafkaConnectionDetails connectionDetails) { + return KafkaConfigBuilder.of(properties).streams().withConnectionDetails(connectionDetails).build(); + } + + @Override + public Security security(KafkaProperties properties) { + return properties.getStreams().getSecurity(); + } + + @Test + void streamsPropertiesAreApplied() { + KafkaProperties properties = new KafkaProperties(); + KafkaProperties.Streams streams = properties.getStreams(); + streams.setApplicationId("app-id"); + streams.setBootstrapServers(List.of("streams:9092")); + streams.setStateStoreCacheMaxSize(DataSize.ofMegabytes(10)); + streams.setClientId("streams-client"); + streams.setReplicationFactor(3); + streams.setStateDir("/tmp/state"); + streams.getProperties().put("streams.custom", "value"); + Map config = build(properties); + assertThat(config).containsEntry("application.id", "app-id") + .containsEntry(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, List.of("streams:9092")) + .containsEntry("statestore.cache.max.bytes", (int) DataSize.ofMegabytes(10).toBytes()) + .containsEntry(StreamsConfig.CLIENT_ID_CONFIG, "streams-client") + .containsEntry("replication.factor", 3) + .containsEntry("state.dir", "/tmp/state") + .containsEntry("streams.custom", "value"); + } + + } + +} diff --git a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaPropertiesTests.java b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaPropertiesTests.java index 3d79ec287da..6d154e33a6e 100644 --- a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaPropertiesTests.java +++ b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaPropertiesTests.java @@ -17,23 +17,18 @@ package org.springframework.boot.kafka.autoconfigure; import java.util.Collections; -import java.util.Map; -import org.apache.kafka.common.config.SslConfigs; import org.junit.jupiter.api.Test; -import org.springframework.boot.context.properties.source.MutuallyExclusiveConfigurationPropertiesException; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Admin; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Cleanup; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.IsolationLevel; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Listener; -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; -import static org.assertj.core.api.Assertions.assertThatExceptionOfType; /** * Tests for {@link KafkaProperties}. @@ -69,87 +64,6 @@ class KafkaPropertiesTests { assertThat(listenerProperties.isMissingTopicsFatal()).isEqualTo(container.isMissingTopicsFatal()); } - @Test - void sslPemConfiguration() { - KafkaProperties properties = new KafkaProperties(); - properties.getSsl().setKeyStoreKey("-----BEGINkey"); - properties.getSsl().setTrustStoreCertificates("-----BEGINtrust"); - properties.getSsl().setKeyStoreCertificateChain("-----BEGINchain"); - Map consumerProperties = properties.buildConsumerProperties(); - assertThat(consumerProperties).containsEntry(SslConfigs.SSL_KEYSTORE_KEY_CONFIG, "-----BEGINkey"); - assertThat(consumerProperties).containsEntry(SslConfigs.SSL_TRUSTSTORE_CERTIFICATES_CONFIG, "-----BEGINtrust"); - assertThat(consumerProperties).containsEntry(SslConfigs.SSL_KEYSTORE_CERTIFICATE_CHAIN_CONFIG, - "-----BEGINchain"); - } - - @Test - void sslPemConfigurationWithEmptyBundle() { - KafkaProperties properties = new KafkaProperties(); - properties.getSsl().setKeyStoreKey("-----BEGINkey"); - properties.getSsl().setTrustStoreCertificates("-----BEGINtrust"); - properties.getSsl().setKeyStoreCertificateChain("-----BEGINchain"); - properties.getSsl().setBundle(""); - Map consumerProperties = properties.buildConsumerProperties(); - assertThat(consumerProperties).containsEntry(SslConfigs.SSL_KEYSTORE_KEY_CONFIG, "-----BEGINkey"); - assertThat(consumerProperties).containsEntry(SslConfigs.SSL_TRUSTSTORE_CERTIFICATES_CONFIG, "-----BEGINtrust"); - assertThat(consumerProperties).containsEntry(SslConfigs.SSL_KEYSTORE_CERTIFICATE_CHAIN_CONFIG, - "-----BEGINchain"); - } - - @Test - void sslPropertiesWhenKeyStoreLocationAndKeySetShouldThrowException() { - KafkaProperties properties = new KafkaProperties(); - properties.getSsl().setKeyStoreKey("-----BEGIN"); - properties.getSsl().setKeyStoreLocation(new ClassPathResource("ksLoc")); - assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) - .isThrownBy(properties::buildConsumerProperties); - } - - @Test - void sslPropertiesWhenTrustStoreLocationAndCertificatesSetShouldThrowException() { - KafkaProperties properties = new KafkaProperties(); - properties.getSsl().setTrustStoreLocation(new ClassPathResource("tsLoc")); - properties.getSsl().setTrustStoreCertificates("-----BEGIN"); - assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) - .isThrownBy(properties::buildConsumerProperties); - } - - @Test - void sslPropertiesWhenKeyStoreLocationAndBundleSetShouldThrowException() { - KafkaProperties properties = new KafkaProperties(); - properties.getSsl().setBundle("myBundle"); - properties.getSsl().setKeyStoreLocation(new ClassPathResource("ksLoc")); - assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) - .isThrownBy(properties::buildConsumerProperties); - } - - @Test - void sslPropertiesWhenKeyStoreKeyAndBundleSetShouldThrowException() { - KafkaProperties properties = new KafkaProperties(); - properties.getSsl().setBundle("myBundle"); - properties.getSsl().setKeyStoreKey("-----BEGIN"); - assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) - .isThrownBy(properties::buildConsumerProperties); - } - - @Test - void sslPropertiesWhenTrustStoreLocationAndBundleSetShouldThrowException() { - KafkaProperties properties = new KafkaProperties(); - properties.getSsl().setBundle("myBundle"); - properties.getSsl().setTrustStoreLocation(new ClassPathResource("tsLoc")); - assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) - .isThrownBy(properties::buildConsumerProperties); - } - - @Test - void sslPropertiesWhenTrustStoreCertificatesAndBundleSetShouldThrowException() { - KafkaProperties properties = new KafkaProperties(); - properties.getSsl().setBundle("myBundle"); - properties.getSsl().setTrustStoreCertificates("-----BEGIN"); - assertThatExceptionOfType(MutuallyExclusiveConfigurationPropertiesException.class) - .isThrownBy(properties::buildConsumerProperties); - } - @Test void cleanupConfigDefaultValuesAreConsistent() { CleanupConfig cleanupConfig = new CleanupConfig();