Add nullability annotations to module/spring-boot-kafka

See gh-46587
This commit is contained in:
Moritz Halbritter
2025-08-04 11:27:41 +02:00
parent d878ed5d14
commit 30e7d1eb80
14 changed files with 265 additions and 253 deletions
@@ -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<Object, Object> recordFilterStrategy;
private @Nullable RecordFilterStrategy<Object, Object> recordFilterStrategy;
private KafkaTemplate<Object, Object> replyTemplate;
private @Nullable KafkaTemplate<Object, Object> replyTemplate;
private KafkaAwareTransactionManager<Object, Object> transactionManager;
private @Nullable KafkaAwareTransactionManager<Object, Object> transactionManager;
private ConsumerAwareRebalanceListener rebalanceListener;
private @Nullable ConsumerAwareRebalanceListener rebalanceListener;
private CommonErrorHandler commonErrorHandler;
private @Nullable CommonErrorHandler commonErrorHandler;
private AfterRollbackProcessor<Object, Object> afterRollbackProcessor;
private @Nullable AfterRollbackProcessor<Object, Object> afterRollbackProcessor;
private RecordInterceptor<Object, Object> recordInterceptor;
private @Nullable RecordInterceptor<Object, Object> recordInterceptor;
private BatchInterceptor<Object, Object> batchInterceptor;
private @Nullable BatchInterceptor<Object, Object> batchInterceptor;
private Function<MessageListenerContainer, String> threadNameSupplier;
private @Nullable Function<MessageListenerContainer, String> 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<Object, Object> recordFilterStrategy) {
void setRecordFilterStrategy(@Nullable RecordFilterStrategy<Object, Object> 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<Object, Object> replyTemplate) {
void setReplyTemplate(@Nullable KafkaTemplate<Object, Object> 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<Object, Object> transactionManager) {
void setTransactionManager(@Nullable KafkaAwareTransactionManager<Object, Object> 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<Object, Object> afterRollbackProcessor) {
void setAfterRollbackProcessor(@Nullable AfterRollbackProcessor<Object, Object> 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<Object, Object> recordInterceptor) {
void setRecordInterceptor(@Nullable RecordInterceptor<Object, Object> 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<Object, Object> batchInterceptor) {
void setBatchInterceptor(@Nullable BatchInterceptor<Object, Object> 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<MessageListenerContainer, String> threadNameSupplier) {
void setThreadNameSupplier(@Nullable Function<MessageListenerContainer, String> 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<Object, Object> 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);
@@ -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<Object, Object> recordFilterStrategy;
private final @Nullable RecordFilterStrategy<Object, Object> recordFilterStrategy;
private final BatchMessageConverter batchMessageConverter;
private final KafkaTemplate<Object, Object> kafkaTemplate;
private final @Nullable KafkaTemplate<Object, Object> kafkaTemplate;
private final KafkaAwareTransactionManager<Object, Object> transactionManager;
private final @Nullable KafkaAwareTransactionManager<Object, Object> transactionManager;
private final ConsumerAwareRebalanceListener rebalanceListener;
private final @Nullable ConsumerAwareRebalanceListener rebalanceListener;
private final CommonErrorHandler commonErrorHandler;
private final @Nullable CommonErrorHandler commonErrorHandler;
private final AfterRollbackProcessor<Object, Object> afterRollbackProcessor;
private final @Nullable AfterRollbackProcessor<Object, Object> afterRollbackProcessor;
private final RecordInterceptor<Object, Object> recordInterceptor;
private final @Nullable RecordInterceptor<Object, Object> recordInterceptor;
private final BatchInterceptor<Object, Object> batchInterceptor;
private final @Nullable BatchInterceptor<Object, Object> batchInterceptor;
private final Function<MessageListenerContainer, String> threadNameSupplier;
private final @Nullable Function<MessageListenerContainer, String> threadNameSupplier;
KafkaAnnotationDrivenConfiguration(KafkaProperties properties,
ObjectProvider<RecordMessageConverter> recordMessageConverter,
@@ -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<String, Object> properties, SslBundle sslBundle) {
static void applySslBundle(Map<String, Object> 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<String, Object> properties, String securityProtocol) {
static void applySecurityProtocol(Map<String, Object> 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);
}
@@ -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<String> bootstrapServers, SslBundle sslBundle, String securityProtocol) {
static Configuration of(List<String> bootstrapServers, @Nullable SslBundle sslBundle,
@Nullable String securityProtocol) {
return new Configuration() {
@Override
public List<String> 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;
}
@@ -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());
@@ -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<String, ?> configs;
private @Nullable Map<String, ?> configs;
private volatile SslBundle sslBundle;
private volatile @Nullable SslBundle sslBundle;
@Override
public void configure(Map<String, ?> 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();
}
@@ -17,4 +17,7 @@
/**
* Auto-configuration for Apache Kafka metrics.
*/
@NullMarked
package org.springframework.boot.kafka.autoconfigure.metrics;
import org.jspecify.annotations.NullMarked;
@@ -17,4 +17,7 @@
/**
* Auto-configuration for Apache Kafka.
*/
@NullMarked
package org.springframework.boot.kafka.autoconfigure;
import org.jspecify.annotations.NullMarked;
@@ -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;
@@ -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();
}
@@ -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();
}
@@ -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();
}
@@ -17,4 +17,7 @@
/**
* Support for testcontainers Kafka service connections.
*/
@NullMarked
package org.springframework.boot.kafka.testcontainers;
import org.jspecify.annotations.NullMarked;