Nullability improvements in KafkaProducerFactory classes
Signed-off-by: Soby Chacko <soby.chacko@broadcom.com>
This commit is contained in:
@@ -144,13 +144,13 @@ public class DefaultKafkaProducerFactory<K, V> extends KafkaResourceFactory
|
||||
|
||||
private TransactionIdSuffixStrategy transactionIdSuffixStrategy = new DefaultTransactionIdSuffixStrategy(0);
|
||||
|
||||
private @Nullable Supplier<Serializer<K>> keySerializerSupplier;
|
||||
private @Nullable Supplier<@Nullable Serializer<K>> keySerializerSupplier;
|
||||
|
||||
private @Nullable Supplier<Serializer<V>> valueSerializerSupplier;
|
||||
private @Nullable Supplier<@Nullable Serializer<V>> valueSerializerSupplier;
|
||||
|
||||
private @Nullable Supplier<Serializer<K>> rawKeySerializerSupplier;
|
||||
private @Nullable Supplier<@Nullable Serializer<K>> rawKeySerializerSupplier;
|
||||
|
||||
private @Nullable Supplier<Serializer<V>> rawValueSerializerSupplier;
|
||||
private @Nullable Supplier<@Nullable Serializer<V>> rawValueSerializerSupplier;
|
||||
|
||||
private Duration physicalCloseTimeout = DEFAULT_PHYSICAL_CLOSE_TIMEOUT;
|
||||
|
||||
@@ -230,8 +230,8 @@ public class DefaultKafkaProducerFactory<K, V> extends KafkaResourceFactory
|
||||
* @since 2.3
|
||||
*/
|
||||
public DefaultKafkaProducerFactory(Map<String, Object> configs,
|
||||
@Nullable Supplier<Serializer<K>> keySerializerSupplier,
|
||||
@Nullable Supplier<Serializer<V>> valueSerializerSupplier) {
|
||||
@Nullable Supplier<@Nullable Serializer<K>> keySerializerSupplier,
|
||||
@Nullable Supplier<@Nullable Serializer<V>> valueSerializerSupplier) {
|
||||
|
||||
this(configs, keySerializerSupplier, valueSerializerSupplier, true);
|
||||
}
|
||||
@@ -252,8 +252,8 @@ public class DefaultKafkaProducerFactory<K, V> extends KafkaResourceFactory
|
||||
* @since 2.8.7
|
||||
*/
|
||||
public DefaultKafkaProducerFactory(Map<String, Object> configs,
|
||||
@Nullable Supplier<Serializer<K>> keySerializerSupplier,
|
||||
@Nullable Supplier<Serializer<V>> valueSerializerSupplier, boolean configureSerializers) {
|
||||
@Nullable Supplier<@Nullable Serializer<K>> keySerializerSupplier,
|
||||
@Nullable Supplier<@Nullable Serializer<V>> valueSerializerSupplier, boolean configureSerializers) {
|
||||
|
||||
this.configs = new ConcurrentHashMap<>(configs);
|
||||
this.configureSerializers = configureSerializers;
|
||||
@@ -269,7 +269,8 @@ public class DefaultKafkaProducerFactory<K, V> extends KafkaResourceFactory
|
||||
}
|
||||
}
|
||||
|
||||
private @Nullable Supplier<Serializer<K>> keySerializerSupplier(@Nullable Supplier<Serializer<K>> keySerializerSupplier) {
|
||||
private @Nullable Supplier<@Nullable Serializer<K>> keySerializerSupplier(
|
||||
@Nullable Supplier<@Nullable Serializer<K>> keySerializerSupplier) {
|
||||
this.rawKeySerializerSupplier = keySerializerSupplier;
|
||||
if (!this.configureSerializers) {
|
||||
return keySerializerSupplier;
|
||||
@@ -285,7 +286,8 @@ public class DefaultKafkaProducerFactory<K, V> extends KafkaResourceFactory
|
||||
};
|
||||
}
|
||||
|
||||
private @Nullable Supplier<Serializer<V>> valueSerializerSupplier(@Nullable Supplier<Serializer<V>> valueSerializerSupplier) {
|
||||
private @Nullable Supplier<@Nullable Serializer<V>> valueSerializerSupplier(
|
||||
@Nullable Supplier<@Nullable Serializer<V>> valueSerializerSupplier) {
|
||||
this.rawValueSerializerSupplier = valueSerializerSupplier;
|
||||
if (!this.configureSerializers) {
|
||||
return valueSerializerSupplier;
|
||||
@@ -341,7 +343,7 @@ public class DefaultKafkaProducerFactory<K, V> extends KafkaResourceFactory
|
||||
* @since 2.8
|
||||
* @see #setConfigureSerializers(boolean)
|
||||
*/
|
||||
public void setKeySerializerSupplier(Supplier<Serializer<K>> keySerializerSupplier) {
|
||||
public void setKeySerializerSupplier(@Nullable Supplier<@Nullable Serializer<K>> keySerializerSupplier) {
|
||||
this.keySerializerSupplier = keySerializerSupplier(keySerializerSupplier);
|
||||
}
|
||||
|
||||
@@ -353,7 +355,7 @@ public class DefaultKafkaProducerFactory<K, V> extends KafkaResourceFactory
|
||||
* @since 2.8
|
||||
* @see #setConfigureSerializers(boolean)
|
||||
*/
|
||||
public void setValueSerializerSupplier(Supplier<Serializer<V>> valueSerializerSupplier) {
|
||||
public void setValueSerializerSupplier(@Nullable Supplier<@Nullable Serializer<V>> valueSerializerSupplier) {
|
||||
this.valueSerializerSupplier = valueSerializerSupplier(valueSerializerSupplier);
|
||||
}
|
||||
|
||||
@@ -469,13 +471,13 @@ public class DefaultKafkaProducerFactory<K, V> extends KafkaResourceFactory
|
||||
|
||||
@Override
|
||||
@Nullable
|
||||
public Supplier<Serializer<K>> getKeySerializerSupplier() {
|
||||
public Supplier<@Nullable Serializer<K>> getKeySerializerSupplier() {
|
||||
return this.rawKeySerializerSupplier;
|
||||
}
|
||||
|
||||
@Override
|
||||
@Nullable
|
||||
public Supplier<Serializer<V>> getValueSerializerSupplier() {
|
||||
public Supplier<@Nullable Serializer<V>> getValueSerializerSupplier() {
|
||||
return this.rawValueSerializerSupplier;
|
||||
}
|
||||
|
||||
@@ -696,7 +698,6 @@ public class DefaultKafkaProducerFactory<K, V> extends KafkaResourceFactory
|
||||
return this.transactionIdPrefix != null;
|
||||
}
|
||||
|
||||
@SuppressWarnings("resource")
|
||||
@Override
|
||||
public void destroy() {
|
||||
CloseSafeProducer<K, V> producerToClose;
|
||||
|
||||
@@ -115,7 +115,7 @@ public interface ProducerFactory<K, V> {
|
||||
* @since 2.5
|
||||
*/
|
||||
@Nullable
|
||||
default Supplier<Serializer<V>> getValueSerializerSupplier() {
|
||||
default Supplier<@Nullable Serializer<V>> getValueSerializerSupplier() {
|
||||
return () -> null;
|
||||
}
|
||||
|
||||
@@ -126,7 +126,7 @@ public interface ProducerFactory<K, V> {
|
||||
* @since 2.5
|
||||
*/
|
||||
@Nullable
|
||||
default Supplier<Serializer<K>> getKeySerializerSupplier() {
|
||||
default Supplier<@Nullable Serializer<K>> getKeySerializerSupplier() {
|
||||
return () -> null;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user