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 3d8bbd0f1..325172945 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 @@ -55,7 +55,6 @@ import org.springframework.core.task.SimpleAsyncTaskExecutor; import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter; import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint; import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; -import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.core.MessageProducer; import org.springframework.messaging.MessageChannel; @@ -230,16 +229,12 @@ public class RabbitMessageChannelBinder protected MessageProducer createConsumerEndpoint(ConsumerDestination consumerDestination, String group, ExtendedConsumerProperties properties) { - DirectChannel convertingBridgeChannel = new DirectChannel(); - convertingBridgeChannel.setBeanFactory(this.getBeanFactory()); - String prefix = properties.getExtension().getPrefix(); String destination = consumerDestination.getName(); String prefixStripped = (StringUtils.isEmpty(prefix) || !destination.startsWith(prefix)) ? destination : destination.substring(prefix.length()); String baseQueueName = StringUtils.hasText(group) ? prefixStripped.substring(0, prefixStripped.indexOf(group)) + group : prefixStripped; - convertingBridgeChannel.setBeanName(baseQueueName + ".bridge"); SimpleMessageListenerContainer listenerContainer = new SimpleMessageListenerContainer( this.connectionFactory); listenerContainer.setAcknowledgeMode(properties.getExtension().getAcknowledgeMode()); diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderCleanerTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderCleanerTests.java index ece3cb541..dd4f67401 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderCleanerTests.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderCleanerTests.java @@ -136,15 +136,20 @@ public class RabbitBinderCleanerTests { URI uri = UriComponentsBuilder.fromUriString("http://localhost:15672/api/queues").pathSegment( "{vhost}", "{queue}") .buildAndExpand("/", queueName).encode().toUri(); - while (n++ < 100) { + + Object consumers = null; + while (n++ < 100 && (consumers == null || consumers.equals(Integer.valueOf(state)))) { @SuppressWarnings("unchecked") Map queueInfo = template.getForObject(uri, Map.class); - if (!queueInfo.get("consumers").equals(Integer.valueOf(state))) { - break; + consumers = queueInfo.get("consumers"); + if (consumers == null || consumers.equals(Integer.valueOf(state))) { + Thread.sleep(100); } - Thread.sleep(100); } - assertThat(n < 100).withFailMessage("Consumer state remained at " + state + " after 10 seconds"); + assertThat(consumers).isNotNull(); + + assertThat(n).withFailMessage("Consumer state remained at " + state + " after 10 seconds") + .isLessThan(100); } });