From c4d3e4720184c0893da251ba51d8820b7fe940b0 Mon Sep 17 00:00:00 2001 From: MichMich <12600416+mmichailidis@users.noreply.github.com> Date: Thu, 7 Nov 2019 18:23:09 +0200 Subject: [PATCH] Ft multiple partition support (#272) * Added support for partitioned multiplex * removed debug line * added tests regarding the multiplex feature for multiple instances --- .../RabbitExchangeQueueProvisioner.java | 26 +++++++++++- .../binder/rabbit/RabbitBinderTests.java | 41 +++++++++++++++++++ 2 files changed, 65 insertions(+), 2 deletions(-) diff --git a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java index 543ac8e5d..336c63b02 100644 --- a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java +++ b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java @@ -16,7 +16,9 @@ package org.springframework.cloud.stream.binder.rabbit.provisioning; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.stream.Stream; @@ -39,6 +41,7 @@ import org.springframework.amqp.core.TopicExchange; import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.core.DeclarationExceptionEvent; import org.springframework.amqp.rabbit.core.RabbitAdmin; +import org.springframework.beans.BeanUtils; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.DefaultListableBeanFactory; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; @@ -63,6 +66,7 @@ import org.springframework.util.StringUtils; * @author Soby Chacko * @author Gary Russell * @author Oleg Zhurakousky + * @author Michael Michailidis */ // @checkstyle:off public class RabbitExchangeQueueProvisioner @@ -170,8 +174,26 @@ public class RabbitExchangeQueueProvisioner else { String[] provisionedDestinations = Stream .of(StringUtils.tokenizeToStringArray(name, ",", true, true)) - .map(destination -> doProvisionConsumerDestination(destination, group, - properties).getName()) + .flatMap(destination -> { + if (properties.isPartitioned() && !ObjectUtils.isEmpty(properties.getInstanceIndexList())) { + List consumerDestinationNames = new ArrayList<>(); + + for (Integer index : properties.getInstanceIndexList()) { + ExtendedConsumerProperties temporaryProperties = + new ExtendedConsumerProperties<>(properties.getExtension()); + BeanUtils.copyProperties(properties, temporaryProperties); + temporaryProperties.setInstanceIndex(index); + consumerDestinationNames.add(doProvisionConsumerDestination(destination, group, + temporaryProperties).getName()); + } + + return consumerDestinationNames.stream(); + } + else { + return Stream.of(doProvisionConsumerDestination(destination, group, + properties).getName()); + } + }) .toArray(String[]::new); consumerDestination = new RabbitConsumerDestination( StringUtils.arrayToCommaDelimitedString(provisionedDestinations), diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java index ef0d46a0a..7c9dae8de 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java @@ -88,6 +88,7 @@ import org.springframework.cloud.stream.binder.rabbit.properties.RabbitProducerP import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchangeQueueProvisioner; import org.springframework.cloud.stream.binder.test.junit.rabbit.RabbitTestSupport; import org.springframework.cloud.stream.config.BindingProperties; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationListener; import org.springframework.context.ConfigurableApplicationContext; @@ -417,6 +418,46 @@ public class RabbitBinderTests extends assertThat(endpoint.isRunning()).isFalse(); } + @Test + public void testMultiplexOnPartitionedConsumer() throws Exception { + final ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + RabbitTestSupport.RabbitProxy proxy = new RabbitTestSupport.RabbitProxy(); + CachingConnectionFactory cf = new CachingConnectionFactory("localhost", + proxy.getPort()); + + final RabbitExchangeQueueProvisioner rabbitExchangeQueueProvisioner = new RabbitExchangeQueueProvisioner(cf); + + consumerProperties.setMultiplex(true); + consumerProperties.setPartitioned(true); + consumerProperties.setInstanceIndexList(Arrays.asList(1, 2, 3)); + + final ConsumerDestination consumerDestination = rabbitExchangeQueueProvisioner.provisionConsumerDestination("foo", "boo", consumerProperties); + + final String name = consumerDestination.getName(); + + assertThat(name).isEqualTo("foo.boo-1,foo.boo-2,foo.boo-3"); + } + + @Test + public void testMultiplexOnPartitionedConsumerWithMultipleDestinations() throws Exception { + final ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + RabbitTestSupport.RabbitProxy proxy = new RabbitTestSupport.RabbitProxy(); + CachingConnectionFactory cf = new CachingConnectionFactory("localhost", + proxy.getPort()); + + final RabbitExchangeQueueProvisioner rabbitExchangeQueueProvisioner = new RabbitExchangeQueueProvisioner(cf); + + consumerProperties.setMultiplex(true); + consumerProperties.setPartitioned(true); + consumerProperties.setInstanceIndexList(Arrays.asList(1, 2, 3)); + + final ConsumerDestination consumerDestination = rabbitExchangeQueueProvisioner.provisionConsumerDestination("foo,qaa", "boo", consumerProperties); + + final String name = consumerDestination.getName(); + + assertThat(name).isEqualTo("foo.boo-1,foo.boo-2,foo.boo-3,qaa.boo-1,qaa.boo-2,qaa.boo-3"); + } + @Test public void testConsumerPropertiesWithUserInfrastructureNoBind() throws Exception { RabbitAdmin admin = new RabbitAdmin(this.rabbitAvailableRule.getResource());