From 41bf4d573fbc7f815fc4dffb3c08a03982897da7 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 11 Nov 2019 11:36:15 -0500 Subject: [PATCH] Configure consumer customizer Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/1841 * Add test --- .../config/RabbitMessageChannelBinderConfiguration.java | 6 +++++- .../rabbit/integration/RabbitBinderModuleTests.java | 9 +++++++++ 2 files changed, 14 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java index fb405fa9f..7ceb3e262 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java @@ -35,12 +35,14 @@ import org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder import org.springframework.cloud.stream.binder.rabbit.properties.RabbitBinderConfigurationProperties; import org.springframework.cloud.stream.binder.rabbit.properties.RabbitExtendedBindingProperties; import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchangeQueueProvisioner; +import org.springframework.cloud.stream.config.ConsumerEndpointCustomizer; 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.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; +import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter; import org.springframework.integration.amqp.inbound.AmqpMessageSource; import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint; import org.springframework.lang.Nullable; @@ -76,7 +78,8 @@ public class RabbitMessageChannelBinderConfiguration { RabbitMessageChannelBinder rabbitMessageChannelBinder( @Nullable ListenerContainerCustomizer listenerContainerCustomizer, @Nullable MessageSourceCustomizer sourceCustomizer, - @Nullable ProducerMessageHandlerCustomizer producerMessageHandlerCustomizer) { + @Nullable ProducerMessageHandlerCustomizer producerMessageHandlerCustomizer, + @Nullable ConsumerEndpointCustomizer consumerCustomizer) { RabbitMessageChannelBinder binder = new RabbitMessageChannelBinder( this.rabbitConnectionFactory, this.rabbitProperties, @@ -88,6 +91,7 @@ public class RabbitMessageChannelBinderConfiguration { binder.setNodes(this.rabbitBinderConfigurationProperties.getNodes()); binder.setExtendedBindingProperties(this.rabbitExtendedBindingProperties); binder.setProducerMessageHandlerCustomizer(producerMessageHandlerCustomizer); + binder.setConsumerEndpointCustomizer(consumerCustomizer); return binder; } diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/integration/RabbitBinderModuleTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/integration/RabbitBinderModuleTests.java index 4f8fc9504..2dc11a73b 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/integration/RabbitBinderModuleTests.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/integration/RabbitBinderModuleTests.java @@ -56,12 +56,14 @@ import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerP import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties; import org.springframework.cloud.stream.binder.test.junit.rabbit.RabbitTestSupport; import org.springframework.cloud.stream.binding.BindingService; +import org.springframework.cloud.stream.config.ConsumerEndpointCustomizer; 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.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; +import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter; import org.springframework.integration.amqp.inbound.AmqpMessageSource; import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint; import org.springframework.integration.channel.DirectChannel; @@ -169,6 +171,8 @@ public class RabbitBinderModuleTests { .getPropertyValue("consumerBindings"); // @checkstyle:on Binding inputBinding = consumerBindings.get("input").get(0); + assertThat(TestUtils.getPropertyValue(inputBinding, "lifecycle.beanName")) + .isEqualTo("setByCustomizer:someGroup"); SimpleMessageListenerContainer container = TestUtils.getPropertyValue( inputBinding, "lifecycle.messageListenerContainer", SimpleMessageListenerContainer.class); @@ -362,6 +366,11 @@ public class RabbitBinderModuleTests { return (handler, destinationName) -> handler.setBeanName("setByCustomizer:" + destinationName); } + @Bean + public ConsumerEndpointCustomizer adapterCustomizer() { + return (producer, dest, grp) -> producer.setBeanName("setByCustomizer:" + grp); + } + } public static class ConnectionFactoryConfiguration {