diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index f8f29abce..a7e8b048f 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -71,6 +71,7 @@ import org.springframework.cloud.stream.binder.rabbit.properties.RabbitExtendedB import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerProperties; import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchangeQueueProvisioner; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; +import org.springframework.cloud.stream.config.MessageSourceCustomizer; import org.springframework.cloud.stream.provisioning.ConsumerDestination; import org.springframework.cloud.stream.provisioning.ProducerDestination; import org.springframework.context.support.GenericApplicationContext; @@ -163,14 +164,25 @@ public class RabbitMessageChannelBinder extends public RabbitMessageChannelBinder(ConnectionFactory connectionFactory, RabbitProperties rabbitProperties, RabbitExchangeQueueProvisioner provisioningProvider) { - this(connectionFactory, rabbitProperties, provisioningProvider, null); + + this(connectionFactory, rabbitProperties, provisioningProvider, null, null); } public RabbitMessageChannelBinder(ConnectionFactory connectionFactory, RabbitProperties rabbitProperties, RabbitExchangeQueueProvisioner provisioningProvider, ListenerContainerCustomizer containerCustomizer) { - super(new String[0], provisioningProvider, containerCustomizer); + + this(connectionFactory, rabbitProperties, provisioningProvider, containerCustomizer, null); + } + + public RabbitMessageChannelBinder(ConnectionFactory connectionFactory, + RabbitProperties rabbitProperties, + RabbitExchangeQueueProvisioner provisioningProvider, + ListenerContainerCustomizer containerCustomizer, + MessageSourceCustomizer sourceCustomizer) { + + super(new String[0], provisioningProvider, containerCustomizer, sourceCustomizer); Assert.notNull(connectionFactory, "connectionFactory must not be null"); Assert.notNull(rabbitProperties, "rabbitProperties must not be null"); this.connectionFactory = connectionFactory; @@ -539,11 +551,13 @@ public class RabbitMessageChannelBinder extends protected PolledConsumerResources createPolledConsumerResources(String name, String group, ConsumerDestination destination, ExtendedConsumerProperties consumerProperties) { + Assert.isTrue(!consumerProperties.isMultiplex(), "The Spring Integration polled MessageSource does not currently support muiltiple queues"); AmqpMessageSource source = new AmqpMessageSource(this.connectionFactory, destination.getName()); source.setRawMessageHeader(true); + getMessageSourceCustomizer().configure(source, destination.getName(), group); return new PolledConsumerResources(source, registerErrorInfrastructure( destination, group, consumerProperties, true)); } 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 bd18afb97..bfbb19ce6 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 @@ -36,9 +36,11 @@ import org.springframework.cloud.stream.binder.rabbit.properties.RabbitBinderCon import org.springframework.cloud.stream.binder.rabbit.properties.RabbitExtendedBindingProperties; import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchangeQueueProvisioner; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; +import org.springframework.cloud.stream.config.MessageSourceCustomizer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; +import org.springframework.integration.amqp.inbound.AmqpMessageSource; import org.springframework.lang.Nullable; /** @@ -70,11 +72,12 @@ public class RabbitMessageChannelBinderConfiguration { @Bean RabbitMessageChannelBinder rabbitMessageChannelBinder( - @Nullable ListenerContainerCustomizer listenerContainerCustomizer) - throws Exception { + @Nullable ListenerContainerCustomizer listenerContainerCustomizer, + @Nullable MessageSourceCustomizer sourceCustomizer) { + RabbitMessageChannelBinder binder = new RabbitMessageChannelBinder( this.rabbitConnectionFactory, this.rabbitProperties, - provisioningProvider(), listenerContainerCustomizer); + provisioningProvider(), listenerContainerCustomizer, sourceCustomizer); binder.setAdminAddresses( this.rabbitBinderConfigurationProperties.getAdminAddresses()); binder.setCompressingPostProcessor(gZipPostProcessor()); 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 3c34a9131..e91623441 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 @@ -1,5 +1,5 @@ /* - * Copyright 2015-2018 the original author or authors. + * Copyright 2015-2019 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -43,21 +43,25 @@ import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.cloud.Cloud; import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.BinderFactory; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; +import org.springframework.cloud.stream.binder.PollableMessageSource; import org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder; import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties; 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.ListenerContainerCustomizer; +import org.springframework.cloud.stream.config.MessageSourceCustomizer; import org.springframework.cloud.stream.messaging.Processor; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; +import org.springframework.integration.amqp.inbound.AmqpMessageSource; import org.springframework.integration.channel.DirectChannel; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.support.GenericMessage; @@ -150,6 +154,7 @@ public class RabbitBinderModuleTests { public void testParentConnectionFactoryInheritedByDefaultAndRabbitSettingsPropagated() { context = new SpringApplicationBuilder(SimpleProcessor.class) .web(WebApplicationType.NONE).run("--server.port=0", + "--spring.cloud.stream.bindings.source.group=someGroup", "--spring.cloud.stream.bindings.input.group=someGroup", "--spring.cloud.stream.rabbit.bindings.input.consumer.transacted=true", "--spring.cloud.stream.rabbit.bindings.output.producer.transacted=true"); @@ -198,6 +203,9 @@ public class RabbitBinderModuleTests { ConnectionNameStrategy cns = TestUtils.getPropertyValue(cf, "connectionNameStrategy", ConnectionNameStrategy.class); assertThat(cns.obtainNewConnectionName(cf)).startsWith("rabbitConnectionFactory"); + assertThat(TestUtils.getPropertyValue(consumerBindings.get("source").get(0), + "target.source.h.advised.targetSource.target.beanName")) + .isEqualTo("setByCustomizer:someGroup"); } @Test @@ -340,7 +348,7 @@ public class RabbitBinderModuleTests { assertThat(rabbitConsumerProperties.getMaxConcurrency()).isEqualTo(4); } - @EnableBinding(Processor.class) + @EnableBinding({ Processor.class, PMS.class }) @SpringBootApplication public static class SimpleProcessor { @@ -350,6 +358,11 @@ public class RabbitBinderModuleTests { "setByCustomizerForQueue:" + q + (g == null ? "" : ",andGroup:" + g)); } + @Bean + public MessageSourceCustomizer sourceCustomizer() { + return (s, q, g) -> s.setBeanName("setByCustomizer:" + g); + } + } public static class ConnectionFactoryConfiguration { @@ -375,4 +388,11 @@ public class RabbitBinderModuleTests { } + public interface PMS { + + @Input + PollableMessageSource source(); + + } + }