diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java index 695200787..d41ef2e98 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java @@ -92,9 +92,9 @@ public class ReactorKafkaBinder private ProducerConfigCustomizer producerConfigCustomizer; - private ReceiverOptionsCustomizer receiverOptionsCustomizer = (name, opts) -> opts; + private ReceiverOptionsCustomizer receiverOptionsCustomizer = (name, opts) -> opts; - private SenderOptionsCustomizer senderOptionsCustomizer = (name, opts) -> opts; + private SenderOptionsCustomizer senderOptionsCustomizer = (name, opts) -> opts; public ReactorKafkaBinder(KafkaBinderConfigurationProperties configurationProperties, KafkaTopicProvisioner provisioner) { @@ -111,6 +111,7 @@ public class ReactorKafkaBinder this.producerConfigCustomizer = producerConfigCustomizer; } + @SuppressWarnings({ "rawtypes", "unchecked" }) public void receiverOptionsCustomizers(ObjectProvider customizers) { if (customizers.getIfUnique() != null) { this.receiverOptionsCustomizer = customizers.getIfUnique(); @@ -120,7 +121,7 @@ public class ReactorKafkaBinder ReceiverOptionsCustomizer customizer = (name, opts) -> { ReceiverOptions last = null; for (ReceiverOptionsCustomizer cust: list) { - last = cust.apply(name, opts); + last = (ReceiverOptions) cust.apply(name, opts); } return last; }; @@ -130,6 +131,7 @@ public class ReactorKafkaBinder } } + @SuppressWarnings({ "rawtypes", "unchecked" }) public void senderOptionsCustomizers(ObjectProvider customizers) { if (customizers.getIfUnique() != null) { this.senderOptionsCustomizer = customizers.getIfUnique(); @@ -139,7 +141,7 @@ public class ReactorKafkaBinder SenderOptionsCustomizer customizer = (name, opts) -> { SenderOptions last = null; for (SenderOptionsCustomizer cust: list) { - last = cust.apply(name, opts); + last = (SenderOptions) cust.apply(name, opts); } return last; }; diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReceiverOptionsCustomizer.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReceiverOptionsCustomizer.java index f00fabe71..2f2fbc230 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReceiverOptionsCustomizer.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReceiverOptionsCustomizer.java @@ -27,13 +27,18 @@ import org.springframework.core.Ordered; * the second is the {@link ReceiverOptions} to customize. Return the customized * {@link ReceiverOptions}. Applied in {@link Ordered order} if multiple customizers are * found. + *

+ * Generic since 4.0.3 to allow customization of deserializers. + * + * @param the key type. + * @param the value type. * * @author Gary Russell * @since 4.0.2 * */ -public interface ReceiverOptionsCustomizer - extends BiFunction, ReceiverOptions>, Ordered { +public interface ReceiverOptionsCustomizer + extends BiFunction, ReceiverOptions>, Ordered { @Override default int getOrder() { diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/SenderOptionsCustomizer.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/SenderOptionsCustomizer.java index 0f43af6a4..37070e865 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/SenderOptionsCustomizer.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/SenderOptionsCustomizer.java @@ -27,13 +27,18 @@ import org.springframework.core.Ordered; * second is the {@link SenderOptions} to customize. Return the customized * {@link SenderOptions}. Applied in {@link Ordered order} if multiple customizers are * found. + *

+ * Generic since 4.0.3 to allow customization of serializers. + * + * @param the key type. + * @param the value type. * * @author Gary Russell * @since 4.0.2 * */ -public interface SenderOptionsCustomizer - extends BiFunction, SenderOptions>, Ordered { +public interface SenderOptionsCustomizer + extends BiFunction, SenderOptions>, Ordered { @Override default int getOrder() { diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderIntegrationTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderIntegrationTests.java index 4db2d0011..6be427eee 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderIntegrationTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderIntegrationTests.java @@ -26,6 +26,7 @@ import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.StringDeserializer; import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; @@ -167,7 +168,7 @@ class ReactorKafkaBinderIntegrationTests { } @Bean - ReceiverOptionsCustomizer cust1() { + ReceiverOptionsCustomizer cust1() { return (t, u) -> { recOptsCustOrder.add("one"); return u; @@ -175,12 +176,13 @@ class ReactorKafkaBinderIntegrationTests { } @Bean - ReceiverOptionsCustomizer cust2() { - return new ReceiverOptionsCustomizer() { + ReceiverOptionsCustomizer cust2() { + return new ReceiverOptionsCustomizer<>() { @Override - public ReceiverOptions apply(String t, ReceiverOptions u) { + public ReceiverOptions apply(String t, ReceiverOptions u) { recOptsCustOrder.add("two"); + u.withKeyDeserializer(new StringDeserializer()); return u; } @@ -189,7 +191,6 @@ class ReactorKafkaBinderIntegrationTests { return -1; } - }; } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderTests.java index 3bff3364d..a1e2a2c69 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderTests.java @@ -220,6 +220,7 @@ public class ReactorKafkaBinderTests { provisioner.setMetadataRetryOperations(new RetryTemplate()); ReactorKafkaBinder binder = new ReactorKafkaBinder(binderProps, provisioner); binder.setApplicationContext(new GenericApplicationContext()); + @SuppressWarnings("rawtypes") ObjectProvider cust = mock(ObjectProvider.class); AtomicBoolean custCalled = new AtomicBoolean(); given(cust.getIfUnique()).willReturn((name, opts) -> {