From 5794fb983c1969d2a4019cbe6dae24844088084b Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 23 Oct 2019 14:05:42 -0400 Subject: [PATCH] Add ProducerMessageHandlerCustomizer support Related to https://github.com/spring-cloud/spring-cloud-stream-binder-rabbit/issues/265 and https://github.com/spring-cloud/spring-cloud-stream/pull/1828 --- .../config/KafkaBinderConfiguration.java | 4 ++++ .../integration/KafkaBinderActuatorTests.java | 19 ++++++++++++++++++- 2 files changed, 22 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java index 1d9a4fb62..df5a4e340 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java @@ -42,12 +42,14 @@ import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProv import org.springframework.cloud.stream.binder.kafka.utils.DlqPartitionFunction; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; import org.springframework.cloud.stream.config.MessageSourceCustomizer; +import org.springframework.cloud.stream.config.ProducerMessageHandlerCustomizer; import org.springframework.context.ApplicationContext; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; import org.springframework.integration.kafka.inbound.KafkaMessageSource; +import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.security.jaas.KafkaJaasLoginModuleInitializer; import org.springframework.kafka.support.LoggingProducerListener; @@ -104,6 +106,7 @@ public class KafkaBinderConfiguration { KafkaTopicProvisioner provisioningProvider, @Nullable ListenerContainerCustomizer> listenerContainerCustomizer, @Nullable MessageSourceCustomizer> sourceCustomizer, + @Nullable ProducerMessageHandlerCustomizer> messageHandlerCustomizer, ObjectProvider rebalanceListener, ObjectProvider dlqPartitionFunction) { @@ -114,6 +117,7 @@ public class KafkaBinderConfiguration { kafkaMessageChannelBinder.setProducerListener(this.producerListener); kafkaMessageChannelBinder .setExtendedBindingProperties(this.kafkaExtendedBindingProperties); + kafkaMessageChannelBinder.setProducerMessageHandlerCustomizer(messageHandlerCustomizer); return kafkaMessageChannelBinder; } diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java index 3ed861b44..0be08139c 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java @@ -41,9 +41,12 @@ import org.springframework.cloud.stream.binder.PollableMessageSource; import org.springframework.cloud.stream.binding.BindingService; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; import org.springframework.cloud.stream.config.MessageSourceCustomizer; +import org.springframework.cloud.stream.config.ProducerMessageHandlerCustomizer; +import org.springframework.cloud.stream.messaging.Processor; import org.springframework.cloud.stream.messaging.Sink; import org.springframework.context.annotation.Bean; import org.springframework.integration.kafka.inbound.KafkaMessageSource; +import org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.test.rule.EmbeddedKafkaRule; @@ -58,6 +61,7 @@ import static org.assertj.core.api.Assertions.assertThat; * @author Oleg Zhurakousky * @author Jon Schneider * @author Gary Russell + * * @since 2.0 */ @RunWith(SpringRunner.class) @@ -126,10 +130,18 @@ public class KafkaBinderActuatorTests { consumerBindings.get("source").get(0)).getPropertyValue( "lifecycle.beanName")) .isEqualTo("setByCustomizer:source"); + + Map> producerBindings = (Map>) channelBindingServiceAccessor + .getPropertyValue("producerBindings"); + + assertThat(new DirectFieldAccessor( + producerBindings.get("output")).getPropertyValue( + "lifecycle.beanName")) + .isEqualTo("setByCustomizer:output"); }); } - @EnableBinding({ Sink.class, PMS.class }) + @EnableBinding({ Processor.class, PMS.class }) @EnableAutoConfiguration public static class KafkaMetricsTestConfig { @@ -143,6 +155,11 @@ public class KafkaBinderActuatorTests { return (s, q, g) -> s.setBeanName("setByCustomizer:" + q); } + @Bean + public ProducerMessageHandlerCustomizer> handlerCustomizer() { + return (handler, destinationName) -> handler.setBeanName("setByCustomizer:" + destinationName); + } + @StreamListener(Sink.INPUT) public void process(@SuppressWarnings("unused") String payload) throws InterruptedException { // Artificial slow listener to emulate consumer lag