Configure consumer customizer
Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/1841 * Add test
This commit is contained in:
committed by
Artem Bilan
parent
4c7c0eb3bc
commit
41bf4d573f
@@ -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<AbstractMessageListenerContainer> listenerContainerCustomizer,
|
||||
@Nullable MessageSourceCustomizer<AmqpMessageSource> sourceCustomizer,
|
||||
@Nullable ProducerMessageHandlerCustomizer<AmqpOutboundEndpoint> producerMessageHandlerCustomizer) {
|
||||
@Nullable ProducerMessageHandlerCustomizer<AmqpOutboundEndpoint> producerMessageHandlerCustomizer,
|
||||
@Nullable ConsumerEndpointCustomizer<AmqpInboundChannelAdapter> 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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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<MessageChannel> 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<AmqpInboundChannelAdapter> adapterCustomizer() {
|
||||
return (producer, dest, grp) -> producer.setBeanName("setByCustomizer:" + grp);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
public static class ConnectionFactoryConfiguration {
|
||||
|
||||
Reference in New Issue
Block a user