diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java index c0fe33e7..b9e9488b 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java @@ -144,13 +144,13 @@ public class DefaultKafkaProducerFactory extends KafkaResourceFactory private TransactionIdSuffixStrategy transactionIdSuffixStrategy = new DefaultTransactionIdSuffixStrategy(0); - private @Nullable Supplier> keySerializerSupplier; + private @Nullable Supplier<@Nullable Serializer> keySerializerSupplier; - private @Nullable Supplier> valueSerializerSupplier; + private @Nullable Supplier<@Nullable Serializer> valueSerializerSupplier; - private @Nullable Supplier> rawKeySerializerSupplier; + private @Nullable Supplier<@Nullable Serializer> rawKeySerializerSupplier; - private @Nullable Supplier> rawValueSerializerSupplier; + private @Nullable Supplier<@Nullable Serializer> rawValueSerializerSupplier; private Duration physicalCloseTimeout = DEFAULT_PHYSICAL_CLOSE_TIMEOUT; @@ -230,8 +230,8 @@ public class DefaultKafkaProducerFactory extends KafkaResourceFactory * @since 2.3 */ public DefaultKafkaProducerFactory(Map configs, - @Nullable Supplier> keySerializerSupplier, - @Nullable Supplier> valueSerializerSupplier) { + @Nullable Supplier<@Nullable Serializer> keySerializerSupplier, + @Nullable Supplier<@Nullable Serializer> valueSerializerSupplier) { this(configs, keySerializerSupplier, valueSerializerSupplier, true); } @@ -252,8 +252,8 @@ public class DefaultKafkaProducerFactory extends KafkaResourceFactory * @since 2.8.7 */ public DefaultKafkaProducerFactory(Map configs, - @Nullable Supplier> keySerializerSupplier, - @Nullable Supplier> valueSerializerSupplier, boolean configureSerializers) { + @Nullable Supplier<@Nullable Serializer> keySerializerSupplier, + @Nullable Supplier<@Nullable Serializer> valueSerializerSupplier, boolean configureSerializers) { this.configs = new ConcurrentHashMap<>(configs); this.configureSerializers = configureSerializers; @@ -269,7 +269,8 @@ public class DefaultKafkaProducerFactory extends KafkaResourceFactory } } - private @Nullable Supplier> keySerializerSupplier(@Nullable Supplier> keySerializerSupplier) { + private @Nullable Supplier<@Nullable Serializer> keySerializerSupplier( + @Nullable Supplier<@Nullable Serializer> keySerializerSupplier) { this.rawKeySerializerSupplier = keySerializerSupplier; if (!this.configureSerializers) { return keySerializerSupplier; @@ -285,7 +286,8 @@ public class DefaultKafkaProducerFactory extends KafkaResourceFactory }; } - private @Nullable Supplier> valueSerializerSupplier(@Nullable Supplier> valueSerializerSupplier) { + private @Nullable Supplier<@Nullable Serializer> valueSerializerSupplier( + @Nullable Supplier<@Nullable Serializer> valueSerializerSupplier) { this.rawValueSerializerSupplier = valueSerializerSupplier; if (!this.configureSerializers) { return valueSerializerSupplier; @@ -341,7 +343,7 @@ public class DefaultKafkaProducerFactory extends KafkaResourceFactory * @since 2.8 * @see #setConfigureSerializers(boolean) */ - public void setKeySerializerSupplier(Supplier> keySerializerSupplier) { + public void setKeySerializerSupplier(@Nullable Supplier<@Nullable Serializer> keySerializerSupplier) { this.keySerializerSupplier = keySerializerSupplier(keySerializerSupplier); } @@ -353,7 +355,7 @@ public class DefaultKafkaProducerFactory extends KafkaResourceFactory * @since 2.8 * @see #setConfigureSerializers(boolean) */ - public void setValueSerializerSupplier(Supplier> valueSerializerSupplier) { + public void setValueSerializerSupplier(@Nullable Supplier<@Nullable Serializer> valueSerializerSupplier) { this.valueSerializerSupplier = valueSerializerSupplier(valueSerializerSupplier); } @@ -469,13 +471,13 @@ public class DefaultKafkaProducerFactory extends KafkaResourceFactory @Override @Nullable - public Supplier> getKeySerializerSupplier() { + public Supplier<@Nullable Serializer> getKeySerializerSupplier() { return this.rawKeySerializerSupplier; } @Override @Nullable - public Supplier> getValueSerializerSupplier() { + public Supplier<@Nullable Serializer> getValueSerializerSupplier() { return this.rawValueSerializerSupplier; } @@ -696,7 +698,6 @@ public class DefaultKafkaProducerFactory extends KafkaResourceFactory return this.transactionIdPrefix != null; } - @SuppressWarnings("resource") @Override public void destroy() { CloseSafeProducer producerToClose; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java index 3183c99f..ae775b0b 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/ProducerFactory.java @@ -115,7 +115,7 @@ public interface ProducerFactory { * @since 2.5 */ @Nullable - default Supplier> getValueSerializerSupplier() { + default Supplier<@Nullable Serializer> getValueSerializerSupplier() { return () -> null; } @@ -126,7 +126,7 @@ public interface ProducerFactory { * @since 2.5 */ @Nullable - default Supplier> getKeySerializerSupplier() { + default Supplier<@Nullable Serializer> getKeySerializerSupplier() { return () -> null; }