From 30e7d1eb80b694667759ffae82a632e44248f26d Mon Sep 17 00:00:00 2001 From: Moritz Halbritter Date: Fri, 1 Aug 2025 14:23:52 +0200 Subject: [PATCH] Add nullability annotations to module/spring-boot-kafka See gh-46587 --- ...fkaListenerContainerFactoryConfigurer.java | 57 +-- .../KafkaAnnotationDrivenConfiguration.java | 22 +- .../autoconfigure/KafkaAutoConfiguration.java | 7 +- .../autoconfigure/KafkaConnectionDetails.java | 17 +- .../kafka/autoconfigure/KafkaProperties.java | 351 +++++++++--------- .../PropertiesKafkaConnectionDetails.java | 12 +- .../SslBundleSslEngineFactory.java | 14 +- .../autoconfigure/metrics/package-info.java | 3 + .../kafka/autoconfigure/package-info.java | 3 + .../metrics/autoconfigure/package-info.java | 20 - ...afkaContainerConnectionDetailsFactory.java | 3 +- ...afkaContainerConnectionDetailsFactory.java | 3 +- ...andaContainerConnectionDetailsFactory.java | 3 +- .../kafka/testcontainers/package-info.java | 3 + 14 files changed, 265 insertions(+), 253 deletions(-) delete mode 100644 module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/metrics/autoconfigure/package-info.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 6171a6ea704..9912ebedd85 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 @@ -19,6 +19,8 @@ package org.springframework.boot.kafka.autoconfigure; import java.time.Duration; import java.util.function.Function; +import org.jspecify.annotations.Nullable; + import org.springframework.boot.context.properties.PropertyMapper; import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Listener; import org.springframework.core.task.SimpleAsyncTaskExecutor; @@ -36,6 +38,7 @@ import org.springframework.kafka.listener.adapter.RecordFilterStrategy; import org.springframework.kafka.support.converter.BatchMessageConverter; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.kafka.transaction.KafkaAwareTransactionManager; +import org.springframework.util.Assert; /** * Configure {@link ConcurrentKafkaListenerContainerFactory} with sensible defaults tuned @@ -53,37 +56,37 @@ import org.springframework.kafka.transaction.KafkaAwareTransactionManager; */ public class ConcurrentKafkaListenerContainerFactoryConfigurer { - private KafkaProperties properties; + private @Nullable KafkaProperties properties; - private BatchMessageConverter batchMessageConverter; + private @Nullable BatchMessageConverter batchMessageConverter; - private RecordMessageConverter recordMessageConverter; + private @Nullable RecordMessageConverter recordMessageConverter; - private RecordFilterStrategy recordFilterStrategy; + private @Nullable RecordFilterStrategy recordFilterStrategy; - private KafkaTemplate replyTemplate; + private @Nullable KafkaTemplate replyTemplate; - private KafkaAwareTransactionManager transactionManager; + private @Nullable KafkaAwareTransactionManager transactionManager; - private ConsumerAwareRebalanceListener rebalanceListener; + private @Nullable ConsumerAwareRebalanceListener rebalanceListener; - private CommonErrorHandler commonErrorHandler; + private @Nullable CommonErrorHandler commonErrorHandler; - private AfterRollbackProcessor afterRollbackProcessor; + private @Nullable AfterRollbackProcessor afterRollbackProcessor; - private RecordInterceptor recordInterceptor; + private @Nullable RecordInterceptor recordInterceptor; - private BatchInterceptor batchInterceptor; + private @Nullable BatchInterceptor batchInterceptor; - private Function threadNameSupplier; + private @Nullable Function threadNameSupplier; - private SimpleAsyncTaskExecutor listenerTaskExecutor; + private @Nullable SimpleAsyncTaskExecutor listenerTaskExecutor; /** * Set the {@link KafkaProperties} to use. * @param properties the properties */ - void setKafkaProperties(KafkaProperties properties) { + void setKafkaProperties(@Nullable KafkaProperties properties) { this.properties = properties; } @@ -91,7 +94,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link BatchMessageConverter} to use. * @param batchMessageConverter the message converter */ - void setBatchMessageConverter(BatchMessageConverter batchMessageConverter) { + void setBatchMessageConverter(@Nullable BatchMessageConverter batchMessageConverter) { this.batchMessageConverter = batchMessageConverter; } @@ -99,7 +102,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link RecordMessageConverter} to use. * @param recordMessageConverter the message converter */ - void setRecordMessageConverter(RecordMessageConverter recordMessageConverter) { + void setRecordMessageConverter(@Nullable RecordMessageConverter recordMessageConverter) { this.recordMessageConverter = recordMessageConverter; } @@ -107,7 +110,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link RecordFilterStrategy} to use to filter incoming records. * @param recordFilterStrategy the record filter strategy */ - void setRecordFilterStrategy(RecordFilterStrategy recordFilterStrategy) { + void setRecordFilterStrategy(@Nullable RecordFilterStrategy recordFilterStrategy) { this.recordFilterStrategy = recordFilterStrategy; } @@ -115,7 +118,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link KafkaTemplate} to use to send replies. * @param replyTemplate the reply template */ - void setReplyTemplate(KafkaTemplate replyTemplate) { + void setReplyTemplate(@Nullable KafkaTemplate replyTemplate) { this.replyTemplate = replyTemplate; } @@ -123,7 +126,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link KafkaAwareTransactionManager} to use. * @param transactionManager the transaction manager */ - void setTransactionManager(KafkaAwareTransactionManager transactionManager) { + void setTransactionManager(@Nullable KafkaAwareTransactionManager transactionManager) { this.transactionManager = transactionManager; } @@ -131,7 +134,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link ConsumerAwareRebalanceListener} to use. * @param rebalanceListener the rebalance listener. */ - void setRebalanceListener(ConsumerAwareRebalanceListener rebalanceListener) { + void setRebalanceListener(@Nullable ConsumerAwareRebalanceListener rebalanceListener) { this.rebalanceListener = rebalanceListener; } @@ -139,7 +142,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link CommonErrorHandler} to use. * @param commonErrorHandler the error handler. */ - public void setCommonErrorHandler(CommonErrorHandler commonErrorHandler) { + public void setCommonErrorHandler(@Nullable CommonErrorHandler commonErrorHandler) { this.commonErrorHandler = commonErrorHandler; } @@ -147,7 +150,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link AfterRollbackProcessor} to use. * @param afterRollbackProcessor the after rollback processor */ - void setAfterRollbackProcessor(AfterRollbackProcessor afterRollbackProcessor) { + void setAfterRollbackProcessor(@Nullable AfterRollbackProcessor afterRollbackProcessor) { this.afterRollbackProcessor = afterRollbackProcessor; } @@ -155,7 +158,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link RecordInterceptor} to use. * @param recordInterceptor the record interceptor. */ - void setRecordInterceptor(RecordInterceptor recordInterceptor) { + void setRecordInterceptor(@Nullable RecordInterceptor recordInterceptor) { this.recordInterceptor = recordInterceptor; } @@ -163,7 +166,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the {@link BatchInterceptor} to use. * @param batchInterceptor the batch interceptor. */ - void setBatchInterceptor(BatchInterceptor batchInterceptor) { + void setBatchInterceptor(@Nullable BatchInterceptor batchInterceptor) { this.batchInterceptor = batchInterceptor; } @@ -171,7 +174,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the thread name supplier to use. * @param threadNameSupplier the thread name supplier to use */ - void setThreadNameSupplier(Function threadNameSupplier) { + void setThreadNameSupplier(@Nullable Function threadNameSupplier) { this.threadNameSupplier = threadNameSupplier; } @@ -179,7 +182,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { * Set the executor for threads that poll the consumer. * @param listenerTaskExecutor task executor */ - void setListenerTaskExecutor(SimpleAsyncTaskExecutor listenerTaskExecutor) { + void setListenerTaskExecutor(@Nullable SimpleAsyncTaskExecutor listenerTaskExecutor) { this.listenerTaskExecutor = listenerTaskExecutor; } @@ -199,6 +202,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { private void configureListenerFactory(ConcurrentKafkaListenerContainerFactory factory) { PropertyMapper map = PropertyMapper.get().alwaysApplyingWhenNonNull(); + Assert.state(this.properties != null, "'properties' must not be null"); Listener properties = this.properties.getListener(); map.from(properties::getConcurrency).to(factory::setConcurrency); map.from(properties::isAutoStartup).to(factory::setAutoStartup); @@ -219,6 +223,7 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer { private void configureContainer(ContainerProperties container) { PropertyMapper map = PropertyMapper.get().alwaysApplyingWhenNonNull(); + Assert.state(this.properties != null, "'properties' must not be null"); Listener properties = this.properties.getListener(); map.from(properties::getAckMode).to(container::setAckMode); map.from(properties::getAsyncAcks).to(container::setAsyncAcks); 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 1cfb79f50d0..4182e4fca22 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 @@ -18,6 +18,8 @@ package org.springframework.boot.kafka.autoconfigure; import java.util.function.Function; +import org.jspecify.annotations.Nullable; + import org.springframework.beans.factory.ObjectProvider; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; @@ -62,27 +64,27 @@ class KafkaAnnotationDrivenConfiguration { private final KafkaProperties properties; - private final RecordMessageConverter recordMessageConverter; + private final @Nullable RecordMessageConverter recordMessageConverter; - private final RecordFilterStrategy recordFilterStrategy; + private final @Nullable RecordFilterStrategy recordFilterStrategy; private final BatchMessageConverter batchMessageConverter; - private final KafkaTemplate kafkaTemplate; + private final @Nullable KafkaTemplate kafkaTemplate; - private final KafkaAwareTransactionManager transactionManager; + private final @Nullable KafkaAwareTransactionManager transactionManager; - private final ConsumerAwareRebalanceListener rebalanceListener; + private final @Nullable ConsumerAwareRebalanceListener rebalanceListener; - private final CommonErrorHandler commonErrorHandler; + private final @Nullable CommonErrorHandler commonErrorHandler; - private final AfterRollbackProcessor afterRollbackProcessor; + private final @Nullable AfterRollbackProcessor afterRollbackProcessor; - private final RecordInterceptor recordInterceptor; + private final @Nullable RecordInterceptor recordInterceptor; - private final BatchInterceptor batchInterceptor; + private final @Nullable BatchInterceptor batchInterceptor; - private final Function threadNameSupplier; + private final @Nullable Function threadNameSupplier; KafkaAnnotationDrivenConfiguration(KafkaProperties properties, ObjectProvider recordMessageConverter, 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 53cedf47ea4..d5947ecd931 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 @@ -24,6 +24,7 @@ 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; import org.springframework.aot.hint.RuntimeHints; @@ -240,14 +241,14 @@ public final class KafkaAutoConfiguration { } } - static void applySslBundle(Map properties, SslBundle sslBundle) { + 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, String securityProtocol) { + static void applySecurityProtocol(Map properties, @Nullable String securityProtocol) { if (StringUtils.hasLength(securityProtocol)) { properties.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol); } @@ -256,7 +257,7 @@ public final class KafkaAutoConfiguration { static class KafkaRuntimeHints implements RuntimeHintsRegistrar { @Override - public void registerHints(RuntimeHints hints, ClassLoader classLoader) { + public void registerHints(RuntimeHints hints, @Nullable ClassLoader classLoader) { hints.reflection().registerType(SslBundleSslEngineFactory.class, MemberCategory.INVOKE_PUBLIC_CONSTRUCTORS); } diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConnectionDetails.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConnectionDetails.java index 8e5a3dfeb88..bbd6ac09790 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConnectionDetails.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaConnectionDetails.java @@ -18,6 +18,8 @@ package org.springframework.boot.kafka.autoconfigure; import java.util.List; +import org.jspecify.annotations.Nullable; + import org.springframework.boot.autoconfigure.service.connection.ConnectionDetails; import org.springframework.boot.ssl.SslBundle; @@ -41,7 +43,7 @@ public interface KafkaConnectionDetails extends ConnectionDetails { * Returns the SSL bundle. * @return the SSL bundle */ - default SslBundle getSslBundle() { + default @Nullable SslBundle getSslBundle() { return null; } @@ -49,7 +51,7 @@ public interface KafkaConnectionDetails extends ConnectionDetails { * Returns the security protocol. * @return the security protocol */ - default String getSecurityProtocol() { + default @Nullable String getSecurityProtocol() { return null; } @@ -117,7 +119,8 @@ public interface KafkaConnectionDetails extends ConnectionDetails { * @param securityProtocol the security protocol * @return the configuration */ - static Configuration of(List bootstrapServers, SslBundle sslBundle, String securityProtocol) { + static Configuration of(List bootstrapServers, @Nullable SslBundle sslBundle, + @Nullable String securityProtocol) { return new Configuration() { @Override public List getBootstrapServers() { @@ -125,12 +128,12 @@ public interface KafkaConnectionDetails extends ConnectionDetails { } @Override - public SslBundle getSslBundle() { + public @Nullable SslBundle getSslBundle() { return sslBundle; } @Override - public String getSecurityProtocol() { + public @Nullable String getSecurityProtocol() { return securityProtocol; } }; @@ -146,7 +149,7 @@ public interface KafkaConnectionDetails extends ConnectionDetails { * Returns the SSL bundle. * @return the SSL bundle */ - default SslBundle getSslBundle() { + default @Nullable SslBundle getSslBundle() { return null; } @@ -154,7 +157,7 @@ public interface KafkaConnectionDetails extends ConnectionDetails { * Returns the security protocol. * @return the security protocol */ - default String getSecurityProtocol() { + default @Nullable String getSecurityProtocol() { return null; } 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 3e4f340b41d..6eb697fdb74 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 @@ -32,6 +32,7 @@ import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.config.SslConfigs; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; +import org.jspecify.annotations.Nullable; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.PropertyMapper; @@ -72,7 +73,7 @@ public class KafkaProperties { /** * ID to pass to the server when making requests. Used for server-side logging. */ - private String clientId; + private @Nullable String clientId; /** * Additional properties, common to producers and consumers, used to configure the @@ -108,11 +109,11 @@ public class KafkaProperties { this.bootstrapServers = bootstrapServers; } - public String getClientId() { + public @Nullable String getClientId() { return this.clientId; } - public void setClientId(String clientId) { + public void setClientId(@Nullable String clientId) { this.clientId = clientId; } @@ -241,52 +242,52 @@ public class KafkaProperties { * Frequency with which the consumer offsets are auto-committed to Kafka if * 'enable.auto.commit' is set to true. */ - private Duration autoCommitInterval; + private @Nullable Duration autoCommitInterval; /** * What to do when there is no initial offset in Kafka or if the current offset no * longer exists on the server. */ - private String autoOffsetReset; + private @Nullable String autoOffsetReset; /** * List of host:port pairs to use for establishing the initial connections to the * Kafka cluster. Overrides the global property, for consumers. */ - private List bootstrapServers; + private @Nullable List bootstrapServers; /** * ID to pass to the server when making requests. Used for server-side logging. */ - private String clientId; + private @Nullable String clientId; /** * Whether the consumer's offset is periodically committed in the background. */ - private Boolean enableAutoCommit; + private @Nullable Boolean enableAutoCommit; /** * Maximum amount of time the server blocks before answering the fetch request if * there isn't sufficient data to immediately satisfy the requirement given by * "fetch-min-size". */ - private Duration fetchMaxWait; + private @Nullable Duration fetchMaxWait; /** * Minimum amount of data the server should return for a fetch request. */ - private DataSize fetchMinSize; + private @Nullable DataSize fetchMinSize; /** * Unique string that identifies the consumer group to which this consumer * belongs. */ - private String groupId; + private @Nullable String groupId; /** * Expected time between heartbeats to the consumer coordinator. */ - private Duration heartbeatInterval; + private @Nullable Duration heartbeatInterval; /** * Isolation level for reading messages that have been written transactionally. @@ -306,13 +307,13 @@ public class KafkaProperties { /** * Maximum number of records returned in a single call to poll(). */ - private Integer maxPollRecords; + private @Nullable Integer maxPollRecords; /** * Maximum delay between invocations of poll() when using consumer group * management. */ - private Duration maxPollInterval; + private @Nullable Duration maxPollInterval; /** * Additional consumer-specific properties used to configure the client. @@ -327,75 +328,75 @@ public class KafkaProperties { return this.security; } - public Duration getAutoCommitInterval() { + public @Nullable Duration getAutoCommitInterval() { return this.autoCommitInterval; } - public void setAutoCommitInterval(Duration autoCommitInterval) { + public void setAutoCommitInterval(@Nullable Duration autoCommitInterval) { this.autoCommitInterval = autoCommitInterval; } - public String getAutoOffsetReset() { + public @Nullable String getAutoOffsetReset() { return this.autoOffsetReset; } - public void setAutoOffsetReset(String autoOffsetReset) { + public void setAutoOffsetReset(@Nullable String autoOffsetReset) { this.autoOffsetReset = autoOffsetReset; } - public List getBootstrapServers() { + public @Nullable List getBootstrapServers() { return this.bootstrapServers; } - public void setBootstrapServers(List bootstrapServers) { + public void setBootstrapServers(@Nullable List bootstrapServers) { this.bootstrapServers = bootstrapServers; } - public String getClientId() { + public @Nullable String getClientId() { return this.clientId; } - public void setClientId(String clientId) { + public void setClientId(@Nullable String clientId) { this.clientId = clientId; } - public Boolean getEnableAutoCommit() { + public @Nullable Boolean getEnableAutoCommit() { return this.enableAutoCommit; } - public void setEnableAutoCommit(Boolean enableAutoCommit) { + public void setEnableAutoCommit(@Nullable Boolean enableAutoCommit) { this.enableAutoCommit = enableAutoCommit; } - public Duration getFetchMaxWait() { + public @Nullable Duration getFetchMaxWait() { return this.fetchMaxWait; } - public void setFetchMaxWait(Duration fetchMaxWait) { + public void setFetchMaxWait(@Nullable Duration fetchMaxWait) { this.fetchMaxWait = fetchMaxWait; } - public DataSize getFetchMinSize() { + public @Nullable DataSize getFetchMinSize() { return this.fetchMinSize; } - public void setFetchMinSize(DataSize fetchMinSize) { + public void setFetchMinSize(@Nullable DataSize fetchMinSize) { this.fetchMinSize = fetchMinSize; } - public String getGroupId() { + public @Nullable String getGroupId() { return this.groupId; } - public void setGroupId(String groupId) { + public void setGroupId(@Nullable String groupId) { this.groupId = groupId; } - public Duration getHeartbeatInterval() { + public @Nullable Duration getHeartbeatInterval() { return this.heartbeatInterval; } - public void setHeartbeatInterval(Duration heartbeatInterval) { + public void setHeartbeatInterval(@Nullable Duration heartbeatInterval) { this.heartbeatInterval = heartbeatInterval; } @@ -423,19 +424,19 @@ public class KafkaProperties { this.valueDeserializer = valueDeserializer; } - public Integer getMaxPollRecords() { + public @Nullable Integer getMaxPollRecords() { return this.maxPollRecords; } - public void setMaxPollRecords(Integer maxPollRecords) { + public void setMaxPollRecords(@Nullable Integer maxPollRecords) { this.maxPollRecords = maxPollRecords; } - public Duration getMaxPollInterval() { + public @Nullable Duration getMaxPollInterval() { return this.maxPollInterval; } - public void setMaxPollInterval(Duration maxPollInterval) { + public void setMaxPollInterval(@Nullable Duration maxPollInterval) { this.maxPollInterval = maxPollInterval; } @@ -486,35 +487,35 @@ public class KafkaProperties { * Number of acknowledgments the producer requires the leader to have received * before considering a request complete. */ - private String acks; + private @Nullable String acks; /** * Default batch size. A small batch size will make batching less common and may * reduce throughput (a batch size of zero disables batching entirely). */ - private DataSize batchSize; + private @Nullable DataSize batchSize; /** * List of host:port pairs to use for establishing the initial connections to the * Kafka cluster. Overrides the global property, for producers. */ - private List bootstrapServers; + private @Nullable List bootstrapServers; /** * Total memory size the producer can use to buffer records waiting to be sent to * the server. */ - private DataSize bufferMemory; + private @Nullable DataSize bufferMemory; /** * ID to pass to the server when making requests. Used for server-side logging. */ - private String clientId; + private @Nullable String clientId; /** * Compression type for all data generated by the producer. */ - private String compressionType; + private @Nullable String compressionType; /** * Serializer class for keys. @@ -529,12 +530,12 @@ public class KafkaProperties { /** * When greater than zero, enables retrying of failed sends. */ - private Integer retries; + private @Nullable Integer retries; /** * When non empty, enables transaction support for producer. */ - private String transactionIdPrefix; + private @Nullable String transactionIdPrefix; /** * Additional producer-specific properties used to configure the client. @@ -549,51 +550,51 @@ public class KafkaProperties { return this.security; } - public String getAcks() { + public @Nullable String getAcks() { return this.acks; } - public void setAcks(String acks) { + public void setAcks(@Nullable String acks) { this.acks = acks; } - public DataSize getBatchSize() { + public @Nullable DataSize getBatchSize() { return this.batchSize; } - public void setBatchSize(DataSize batchSize) { + public void setBatchSize(@Nullable DataSize batchSize) { this.batchSize = batchSize; } - public List getBootstrapServers() { + public @Nullable List getBootstrapServers() { return this.bootstrapServers; } - public void setBootstrapServers(List bootstrapServers) { + public void setBootstrapServers(@Nullable List bootstrapServers) { this.bootstrapServers = bootstrapServers; } - public DataSize getBufferMemory() { + public @Nullable DataSize getBufferMemory() { return this.bufferMemory; } - public void setBufferMemory(DataSize bufferMemory) { + public void setBufferMemory(@Nullable DataSize bufferMemory) { this.bufferMemory = bufferMemory; } - public String getClientId() { + public @Nullable String getClientId() { return this.clientId; } - public void setClientId(String clientId) { + public void setClientId(@Nullable String clientId) { this.clientId = clientId; } - public String getCompressionType() { + public @Nullable String getCompressionType() { return this.compressionType; } - public void setCompressionType(String compressionType) { + public void setCompressionType(@Nullable String compressionType) { this.compressionType = compressionType; } @@ -613,19 +614,19 @@ public class KafkaProperties { this.valueSerializer = valueSerializer; } - public Integer getRetries() { + public @Nullable Integer getRetries() { return this.retries; } - public void setRetries(Integer retries) { + public void setRetries(@Nullable Integer retries) { this.retries = retries; } - public String getTransactionIdPrefix() { + public @Nullable String getTransactionIdPrefix() { return this.transactionIdPrefix; } - public void setTransactionIdPrefix(String transactionIdPrefix) { + public void setTransactionIdPrefix(@Nullable String transactionIdPrefix) { this.transactionIdPrefix = transactionIdPrefix; } @@ -661,7 +662,7 @@ public class KafkaProperties { /** * ID to pass to the server when making requests. Used for server-side logging. */ - private String clientId; + private @Nullable String clientId; /** * Additional admin-specific properties used to configure the client. @@ -671,12 +672,12 @@ public class KafkaProperties { /** * Close timeout. */ - private Duration closeTimeout; + private @Nullable Duration closeTimeout; /** * Operation timeout. */ - private Duration operationTimeout; + private @Nullable Duration operationTimeout; /** * Whether to fail fast if the broker is not available on startup. @@ -702,27 +703,27 @@ public class KafkaProperties { return this.security; } - public String getClientId() { + public @Nullable String getClientId() { return this.clientId; } - public void setClientId(String clientId) { + public void setClientId(@Nullable String clientId) { this.clientId = clientId; } - public Duration getCloseTimeout() { + public @Nullable Duration getCloseTimeout() { return this.closeTimeout; } - public void setCloseTimeout(Duration closeTimeout) { + public void setCloseTimeout(@Nullable Duration closeTimeout) { this.closeTimeout = closeTimeout; } - public Duration getOperationTimeout() { + public @Nullable Duration getOperationTimeout() { return this.operationTimeout; } - public void setOperationTimeout(Duration operationTimeout) { + public void setOperationTimeout(@Nullable Duration operationTimeout) { this.operationTimeout = operationTimeout; } @@ -777,7 +778,7 @@ public class KafkaProperties { /** * Kafka streams application.id property; default spring.application.name. */ - private String applicationId; + private @Nullable String applicationId; /** * Whether to auto-start the streams factory bean. @@ -788,28 +789,28 @@ public class KafkaProperties { * List of host:port pairs to use for establishing the initial connections to the * Kafka cluster. Overrides the global property, for streams. */ - private List bootstrapServers; + private @Nullable List bootstrapServers; /** * Maximum size of the in-memory state store cache across all threads. */ - private DataSize stateStoreCacheMaxSize; + private @Nullable DataSize stateStoreCacheMaxSize; /** * ID to pass to the server when making requests. Used for server-side logging. */ - private String clientId; + private @Nullable String clientId; /** * The replication factor for change log topics and repartition topics created by * the stream processing application. */ - private Integer replicationFactor; + private @Nullable Integer replicationFactor; /** * Directory location for the state store. */ - private String stateDir; + private @Nullable String stateDir; /** * Additional Kafka properties used to configure the streams. @@ -828,11 +829,11 @@ public class KafkaProperties { return this.cleanup; } - public String getApplicationId() { + public @Nullable String getApplicationId() { return this.applicationId; } - public void setApplicationId(String applicationId) { + public void setApplicationId(@Nullable String applicationId) { this.applicationId = applicationId; } @@ -844,43 +845,43 @@ public class KafkaProperties { this.autoStartup = autoStartup; } - public List getBootstrapServers() { + public @Nullable List getBootstrapServers() { return this.bootstrapServers; } - public void setBootstrapServers(List bootstrapServers) { + public void setBootstrapServers(@Nullable List bootstrapServers) { this.bootstrapServers = bootstrapServers; } - public DataSize getStateStoreCacheMaxSize() { + public @Nullable DataSize getStateStoreCacheMaxSize() { return this.stateStoreCacheMaxSize; } - public void setStateStoreCacheMaxSize(DataSize stateStoreCacheMaxSize) { + public void setStateStoreCacheMaxSize(@Nullable DataSize stateStoreCacheMaxSize) { this.stateStoreCacheMaxSize = stateStoreCacheMaxSize; } - public String getClientId() { + public @Nullable String getClientId() { return this.clientId; } - public void setClientId(String clientId) { + public void setClientId(@Nullable String clientId) { this.clientId = clientId; } - public Integer getReplicationFactor() { + public @Nullable Integer getReplicationFactor() { return this.replicationFactor; } - public void setReplicationFactor(Integer replicationFactor) { + public void setReplicationFactor(@Nullable Integer replicationFactor) { this.replicationFactor = replicationFactor; } - public String getStateDir() { + public @Nullable String getStateDir() { return this.stateDir; } - public void setStateDir(String stateDir) { + public void setStateDir(@Nullable String stateDir) { this.stateDir = stateDir; } @@ -909,32 +910,32 @@ public class KafkaProperties { /** * Default topic to which messages are sent. */ - private String defaultTopic; + private @Nullable String defaultTopic; /** * Transaction id prefix, override the transaction id prefix in the producer * factory. */ - private String transactionIdPrefix; + private @Nullable String transactionIdPrefix; /** * Whether to enable observation. */ private boolean observationEnabled; - public String getDefaultTopic() { + public @Nullable String getDefaultTopic() { return this.defaultTopic; } - public void setDefaultTopic(String defaultTopic) { + public void setDefaultTopic(@Nullable String defaultTopic) { this.defaultTopic = defaultTopic; } - public String getTransactionIdPrefix() { + public @Nullable String getTransactionIdPrefix() { return this.transactionIdPrefix; } - public void setTransactionIdPrefix(String transactionIdPrefix) { + public void setTransactionIdPrefix(@Nullable String transactionIdPrefix) { this.transactionIdPrefix = transactionIdPrefix; } @@ -972,45 +973,45 @@ public class KafkaProperties { /** * Listener AckMode. See the spring-kafka documentation. */ - private AckMode ackMode; + private @Nullable AckMode ackMode; /** * Support for asynchronous record acknowledgements. Only applies when * spring.kafka.listener.ack-mode is manual or manual-immediate. */ - private Boolean asyncAcks; + private @Nullable Boolean asyncAcks; /** * Prefix for the listener's consumer client.id property. */ - private String clientId; + private @Nullable String clientId; /** * Number of threads to run in the listener containers. */ - private Integer concurrency; + private @Nullable Integer concurrency; /** * Timeout to use when polling the consumer. */ - private Duration pollTimeout; + private @Nullable Duration pollTimeout; /** * Multiplier applied to "pollTimeout" to determine if a consumer is * non-responsive. */ - private Float noPollThreshold; + private @Nullable Float noPollThreshold; /** * Number of records between offset commits when ackMode is "COUNT" or * "COUNT_TIME". */ - private Integer ackCount; + private @Nullable Integer ackCount; /** * Time between offset commits when ackMode is "TIME" or "COUNT_TIME". */ - private Duration ackTime; + private @Nullable Duration ackTime; /** * Sleep interval between Consumer.poll(Duration) calls. @@ -1020,25 +1021,25 @@ public class KafkaProperties { /** * Time between publishing idle consumer events (no data received). */ - private Duration idleEventInterval; + private @Nullable Duration idleEventInterval; /** * Time between publishing idle partition consumer events (no data received for * partition). */ - private Duration idlePartitionEventInterval; + private @Nullable Duration idlePartitionEventInterval; /** * Time between checks for non-responsive consumers. If a duration suffix is not * specified, seconds will be used. */ @DurationUnit(ChronoUnit.SECONDS) - private Duration monitorInterval; + private @Nullable Duration monitorInterval; /** * Whether to log the container configuration during initialization (INFO level). */ - private Boolean logContainerConfig; + private @Nullable Boolean logContainerConfig; /** * Whether the container should fail to start if at least one of the configured @@ -1061,7 +1062,7 @@ public class KafkaProperties { * Whether to instruct the container to change the consumer thread name during * initialization. */ - private Boolean changeConsumerThreadName; + private @Nullable Boolean changeConsumerThreadName; /** * Whether to enable observation. @@ -1071,7 +1072,7 @@ public class KafkaProperties { /** * Time between retries after authentication exceptions. */ - private Duration authExceptionRetryInterval; + private @Nullable Duration authExceptionRetryInterval; public Type getType() { return this.type; @@ -1081,67 +1082,67 @@ public class KafkaProperties { this.type = type; } - public AckMode getAckMode() { + public @Nullable AckMode getAckMode() { return this.ackMode; } - public void setAckMode(AckMode ackMode) { + public void setAckMode(@Nullable AckMode ackMode) { this.ackMode = ackMode; } - public Boolean getAsyncAcks() { + public @Nullable Boolean getAsyncAcks() { return this.asyncAcks; } - public void setAsyncAcks(Boolean asyncAcks) { + public void setAsyncAcks(@Nullable Boolean asyncAcks) { this.asyncAcks = asyncAcks; } - public String getClientId() { + public @Nullable String getClientId() { return this.clientId; } - public void setClientId(String clientId) { + public void setClientId(@Nullable String clientId) { this.clientId = clientId; } - public Integer getConcurrency() { + public @Nullable Integer getConcurrency() { return this.concurrency; } - public void setConcurrency(Integer concurrency) { + public void setConcurrency(@Nullable Integer concurrency) { this.concurrency = concurrency; } - public Duration getPollTimeout() { + public @Nullable Duration getPollTimeout() { return this.pollTimeout; } - public void setPollTimeout(Duration pollTimeout) { + public void setPollTimeout(@Nullable Duration pollTimeout) { this.pollTimeout = pollTimeout; } - public Float getNoPollThreshold() { + public @Nullable Float getNoPollThreshold() { return this.noPollThreshold; } - public void setNoPollThreshold(Float noPollThreshold) { + public void setNoPollThreshold(@Nullable Float noPollThreshold) { this.noPollThreshold = noPollThreshold; } - public Integer getAckCount() { + public @Nullable Integer getAckCount() { return this.ackCount; } - public void setAckCount(Integer ackCount) { + public void setAckCount(@Nullable Integer ackCount) { this.ackCount = ackCount; } - public Duration getAckTime() { + public @Nullable Duration getAckTime() { return this.ackTime; } - public void setAckTime(Duration ackTime) { + public void setAckTime(@Nullable Duration ackTime) { this.ackTime = ackTime; } @@ -1153,35 +1154,35 @@ public class KafkaProperties { this.idleBetweenPolls = idleBetweenPolls; } - public Duration getIdleEventInterval() { + public @Nullable Duration getIdleEventInterval() { return this.idleEventInterval; } - public void setIdleEventInterval(Duration idleEventInterval) { + public void setIdleEventInterval(@Nullable Duration idleEventInterval) { this.idleEventInterval = idleEventInterval; } - public Duration getIdlePartitionEventInterval() { + public @Nullable Duration getIdlePartitionEventInterval() { return this.idlePartitionEventInterval; } - public void setIdlePartitionEventInterval(Duration idlePartitionEventInterval) { + public void setIdlePartitionEventInterval(@Nullable Duration idlePartitionEventInterval) { this.idlePartitionEventInterval = idlePartitionEventInterval; } - public Duration getMonitorInterval() { + public @Nullable Duration getMonitorInterval() { return this.monitorInterval; } - public void setMonitorInterval(Duration monitorInterval) { + public void setMonitorInterval(@Nullable Duration monitorInterval) { this.monitorInterval = monitorInterval; } - public Boolean getLogContainerConfig() { + public @Nullable Boolean getLogContainerConfig() { return this.logContainerConfig; } - public void setLogContainerConfig(Boolean logContainerConfig) { + public void setLogContainerConfig(@Nullable Boolean logContainerConfig) { this.logContainerConfig = logContainerConfig; } @@ -1209,11 +1210,11 @@ public class KafkaProperties { this.autoStartup = autoStartup; } - public Boolean getChangeConsumerThreadName() { + public @Nullable Boolean getChangeConsumerThreadName() { return this.changeConsumerThreadName; } - public void setChangeConsumerThreadName(Boolean changeConsumerThreadName) { + public void setChangeConsumerThreadName(@Nullable Boolean changeConsumerThreadName) { this.changeConsumerThreadName = changeConsumerThreadName; } @@ -1225,11 +1226,11 @@ public class KafkaProperties { this.observationEnabled = observationEnabled; } - public Duration getAuthExceptionRetryInterval() { + public @Nullable Duration getAuthExceptionRetryInterval() { return this.authExceptionRetryInterval; } - public void setAuthExceptionRetryInterval(Duration authExceptionRetryInterval) { + public void setAuthExceptionRetryInterval(@Nullable Duration authExceptionRetryInterval) { this.authExceptionRetryInterval = authExceptionRetryInterval; } @@ -1240,156 +1241,156 @@ public class KafkaProperties { /** * Name of the SSL bundle to use. */ - private String bundle; + private @Nullable String bundle; /** * Password of the private key in either key store key or key store file. */ - private String keyPassword; + private @Nullable String keyPassword; /** * Certificate chain in PEM format with a list of X.509 certificates. */ - private String keyStoreCertificateChain; + private @Nullable String keyStoreCertificateChain; /** * Private key in PEM format with PKCS#8 keys. */ - private String keyStoreKey; + private @Nullable String keyStoreKey; /** * Location of the key store file. */ - private Resource keyStoreLocation; + private @Nullable Resource keyStoreLocation; /** * Store password for the key store file. */ - private String keyStorePassword; + private @Nullable String keyStorePassword; /** * Type of the key store. */ - private String keyStoreType; + private @Nullable String keyStoreType; /** * Trusted certificates in PEM format with X.509 certificates. */ - private String trustStoreCertificates; + private @Nullable String trustStoreCertificates; /** * Location of the trust store file. */ - private Resource trustStoreLocation; + private @Nullable Resource trustStoreLocation; /** * Store password for the trust store file. */ - private String trustStorePassword; + private @Nullable String trustStorePassword; /** * Type of the trust store. */ - private String trustStoreType; + private @Nullable String trustStoreType; /** * SSL protocol to use. */ - private String protocol; + private @Nullable String protocol; - public String getBundle() { + public @Nullable String getBundle() { return this.bundle; } - public void setBundle(String bundle) { + public void setBundle(@Nullable String bundle) { this.bundle = bundle; } - public String getKeyPassword() { + public @Nullable String getKeyPassword() { return this.keyPassword; } - public void setKeyPassword(String keyPassword) { + public void setKeyPassword(@Nullable String keyPassword) { this.keyPassword = keyPassword; } - public String getKeyStoreCertificateChain() { + public @Nullable String getKeyStoreCertificateChain() { return this.keyStoreCertificateChain; } - public void setKeyStoreCertificateChain(String keyStoreCertificateChain) { + public void setKeyStoreCertificateChain(@Nullable String keyStoreCertificateChain) { this.keyStoreCertificateChain = keyStoreCertificateChain; } - public String getKeyStoreKey() { + public @Nullable String getKeyStoreKey() { return this.keyStoreKey; } - public void setKeyStoreKey(String keyStoreKey) { + public void setKeyStoreKey(@Nullable String keyStoreKey) { this.keyStoreKey = keyStoreKey; } - public Resource getKeyStoreLocation() { + public @Nullable Resource getKeyStoreLocation() { return this.keyStoreLocation; } - public void setKeyStoreLocation(Resource keyStoreLocation) { + public void setKeyStoreLocation(@Nullable Resource keyStoreLocation) { this.keyStoreLocation = keyStoreLocation; } - public String getKeyStorePassword() { + public @Nullable String getKeyStorePassword() { return this.keyStorePassword; } - public void setKeyStorePassword(String keyStorePassword) { + public void setKeyStorePassword(@Nullable String keyStorePassword) { this.keyStorePassword = keyStorePassword; } - public String getKeyStoreType() { + public @Nullable String getKeyStoreType() { return this.keyStoreType; } - public void setKeyStoreType(String keyStoreType) { + public void setKeyStoreType(@Nullable String keyStoreType) { this.keyStoreType = keyStoreType; } - public String getTrustStoreCertificates() { + public @Nullable String getTrustStoreCertificates() { return this.trustStoreCertificates; } - public void setTrustStoreCertificates(String trustStoreCertificates) { + public void setTrustStoreCertificates(@Nullable String trustStoreCertificates) { this.trustStoreCertificates = trustStoreCertificates; } - public Resource getTrustStoreLocation() { + public @Nullable Resource getTrustStoreLocation() { return this.trustStoreLocation; } - public void setTrustStoreLocation(Resource trustStoreLocation) { + public void setTrustStoreLocation(@Nullable Resource trustStoreLocation) { this.trustStoreLocation = trustStoreLocation; } - public String getTrustStorePassword() { + public @Nullable String getTrustStorePassword() { return this.trustStorePassword; } - public void setTrustStorePassword(String trustStorePassword) { + public void setTrustStorePassword(@Nullable String trustStorePassword) { this.trustStorePassword = trustStorePassword; } - public String getTrustStoreType() { + public @Nullable String getTrustStoreType() { return this.trustStoreType; } - public void setTrustStoreType(String trustStoreType) { + public void setTrustStoreType(@Nullable String trustStoreType) { this.trustStoreType = trustStoreType; } - public String getProtocol() { + public @Nullable String getProtocol() { return this.protocol; } - public void setProtocol(String protocol) { + public void setProtocol(@Nullable String protocol) { this.protocol = protocol; } @@ -1447,7 +1448,7 @@ public class KafkaProperties { }, this::hasValue); } - private boolean hasValue(Object value) { + private boolean hasValue(@Nullable Object value) { return (value instanceof String string) ? StringUtils.hasText(string) : value != null; } @@ -1525,13 +1526,13 @@ public class KafkaProperties { /** * Security protocol used to communicate with brokers. */ - private String protocol; + private @Nullable String protocol; - public String getProtocol() { + public @Nullable String getProtocol() { return this.protocol; } - public void setProtocol(String protocol) { + public void setProtocol(@Nullable String protocol) { this.protocol = protocol; } diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/PropertiesKafkaConnectionDetails.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/PropertiesKafkaConnectionDetails.java index da25d7bc261..23e48433a84 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/PropertiesKafkaConnectionDetails.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/PropertiesKafkaConnectionDetails.java @@ -18,6 +18,8 @@ package org.springframework.boot.kafka.autoconfigure; import java.util.List; +import org.jspecify.annotations.Nullable; + import org.springframework.boot.kafka.autoconfigure.KafkaProperties.Ssl; import org.springframework.boot.ssl.SslBundle; import org.springframework.boot.ssl.SslBundles; @@ -35,9 +37,9 @@ class PropertiesKafkaConnectionDetails implements KafkaConnectionDetails { private final KafkaProperties properties; - private final SslBundles sslBundles; + private final @Nullable SslBundles sslBundles; - PropertiesKafkaConnectionDetails(KafkaProperties properties, SslBundles sslBundles) { + PropertiesKafkaConnectionDetails(KafkaProperties properties, @Nullable SslBundles sslBundles) { this.properties = properties; this.sslBundles = sslBundles; } @@ -86,16 +88,16 @@ class PropertiesKafkaConnectionDetails implements KafkaConnectionDetails { } @Override - public SslBundle getSslBundle() { + public @Nullable SslBundle getSslBundle() { return getBundle(this.properties.getSsl()); } @Override - public String getSecurityProtocol() { + public @Nullable String getSecurityProtocol() { return this.properties.getSecurity().getProtocol(); } - private SslBundle getBundle(Ssl ssl) { + private @Nullable SslBundle getBundle(Ssl ssl) { if (StringUtils.hasLength(ssl.getBundle())) { Assert.notNull(this.sslBundles, "SSL bundle name has been set but no SSL bundles found in context"); return this.sslBundles.getBundle(ssl.getBundle()); diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/SslBundleSslEngineFactory.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/SslBundleSslEngineFactory.java index 8d227ee3d7f..41241fbabbd 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/SslBundleSslEngineFactory.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/SslBundleSslEngineFactory.java @@ -25,8 +25,10 @@ import javax.net.ssl.SSLEngine; import javax.net.ssl.SSLParameters; import org.apache.kafka.common.security.auth.SslEngineFactory; +import org.jspecify.annotations.Nullable; import org.springframework.boot.ssl.SslBundle; +import org.springframework.util.Assert; /** * An {@link SslEngineFactory} that configures creates an {@link SSLEngine} from an @@ -40,9 +42,9 @@ public class SslBundleSslEngineFactory implements SslEngineFactory { private static final String SSL_BUNDLE_CONFIG_NAME = SslBundle.class.getName(); - private Map configs; + private @Nullable Map configs; - private volatile SslBundle sslBundle; + private volatile @Nullable SslBundle sslBundle; @Override public void configure(Map configs) { @@ -57,6 +59,7 @@ public class SslBundleSslEngineFactory implements SslEngineFactory { @Override public SSLEngine createClientSslEngine(String peerHost, int peerPort, String endpointIdentification) { + Assert.state(this.sslBundle != null, "'sslBundle' must not be null"); SSLEngine sslEngine = this.sslBundle.createSslContext().createSSLEngine(peerHost, peerPort); sslEngine.setUseClientMode(true); SSLParameters sslParams = sslEngine.getSSLParameters(); @@ -67,6 +70,7 @@ public class SslBundleSslEngineFactory implements SslEngineFactory { @Override public SSLEngine createServerSslEngine(String peerHost, int peerPort) { + Assert.state(this.sslBundle != null, "'sslBundle' must not be null"); SSLEngine sslEngine = this.sslBundle.createSslContext().createSSLEngine(peerHost, peerPort); sslEngine.setUseClientMode(false); return sslEngine; @@ -83,12 +87,14 @@ public class SslBundleSslEngineFactory implements SslEngineFactory { } @Override - public KeyStore keystore() { + public @Nullable KeyStore keystore() { + Assert.state(this.sslBundle != null, "'sslBundle' must not be null"); return this.sslBundle.getStores().getKeyStore(); } @Override - public KeyStore truststore() { + public @Nullable KeyStore truststore() { + Assert.state(this.sslBundle != null, "'sslBundle' must not be null"); return this.sslBundle.getStores().getTrustStore(); } diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/metrics/package-info.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/metrics/package-info.java index f6a476f5188..053b35132bf 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/metrics/package-info.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/metrics/package-info.java @@ -17,4 +17,7 @@ /** * Auto-configuration for Apache Kafka metrics. */ +@NullMarked package org.springframework.boot.kafka.autoconfigure.metrics; + +import org.jspecify.annotations.NullMarked; diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/package-info.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/package-info.java index e3d109b4314..bafec787280 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/package-info.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/package-info.java @@ -17,4 +17,7 @@ /** * Auto-configuration for Apache Kafka. */ +@NullMarked package org.springframework.boot.kafka.autoconfigure; + +import org.jspecify.annotations.NullMarked; diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/metrics/autoconfigure/package-info.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/metrics/autoconfigure/package-info.java deleted file mode 100644 index d17e7a667b7..00000000000 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/metrics/autoconfigure/package-info.java +++ /dev/null @@ -1,20 +0,0 @@ -/* - * 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. - */ - -/** - * Auto-configuration for Apache Kafka metrics. - */ -package org.springframework.boot.kafka.metrics.autoconfigure; diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/ApacheKafkaContainerConnectionDetailsFactory.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/ApacheKafkaContainerConnectionDetailsFactory.java index 6dbe317b0f3..4f85a57e320 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/ApacheKafkaContainerConnectionDetailsFactory.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/ApacheKafkaContainerConnectionDetailsFactory.java @@ -18,6 +18,7 @@ package org.springframework.boot.kafka.testcontainers; import java.util.List; +import org.jspecify.annotations.Nullable; import org.testcontainers.kafka.KafkaContainer; import org.springframework.boot.kafka.autoconfigure.KafkaConnectionDetails; @@ -59,7 +60,7 @@ class ApacheKafkaContainerConnectionDetailsFactory } @Override - public SslBundle getSslBundle() { + public @Nullable SslBundle getSslBundle() { return super.getSslBundle(); } diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/ConfluentKafkaContainerConnectionDetailsFactory.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/ConfluentKafkaContainerConnectionDetailsFactory.java index 6e9b767e960..47b8049d2c9 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/ConfluentKafkaContainerConnectionDetailsFactory.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/ConfluentKafkaContainerConnectionDetailsFactory.java @@ -18,6 +18,7 @@ package org.springframework.boot.kafka.testcontainers; import java.util.List; +import org.jspecify.annotations.Nullable; import org.testcontainers.kafka.ConfluentKafkaContainer; import org.springframework.boot.kafka.autoconfigure.KafkaConnectionDetails; @@ -60,7 +61,7 @@ class ConfluentKafkaContainerConnectionDetailsFactory } @Override - public SslBundle getSslBundle() { + public @Nullable SslBundle getSslBundle() { return super.getSslBundle(); } diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/RedpandaContainerConnectionDetailsFactory.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/RedpandaContainerConnectionDetailsFactory.java index 289aa77d780..554c164b8b9 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/RedpandaContainerConnectionDetailsFactory.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/RedpandaContainerConnectionDetailsFactory.java @@ -18,6 +18,7 @@ package org.springframework.boot.kafka.testcontainers; import java.util.List; +import org.jspecify.annotations.Nullable; import org.testcontainers.redpanda.RedpandaContainer; import org.springframework.boot.kafka.autoconfigure.KafkaConnectionDetails; @@ -57,7 +58,7 @@ class RedpandaContainerConnectionDetailsFactory } @Override - public SslBundle getSslBundle() { + public @Nullable SslBundle getSslBundle() { return super.getSslBundle(); } diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/package-info.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/package-info.java index 7c01ce0890a..8bb058fa81b 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/package-info.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/testcontainers/package-info.java @@ -17,4 +17,7 @@ /** * Support for testcontainers Kafka service connections. */ +@NullMarked package org.springframework.boot.kafka.testcontainers; + +import org.jspecify.annotations.NullMarked;