From fa035296f2338d024fcf5f1178d26e69d3a4c360 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 17 Mar 2016 16:11:27 -0400 Subject: [PATCH] GH-439: Add 'transacted' to Rabbit Producer Props Fixes #439 Resolves #401 --- .../rabbit/RabbitMessageChannelBinder.java | 149 ++++++++---------- .../rabbit/RabbitProducerProperties.java | 11 ++ .../binder/rabbit/RabbitBinderTests.java | 3 + 3 files changed, 83 insertions(+), 80 deletions(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index cf1e4d0f1..72ccfafe4 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -86,7 +86,6 @@ import org.springframework.messaging.SubscribableChannel; import org.springframework.retry.interceptor.RetryOperationsInterceptor; import org.springframework.scheduling.TaskScheduler; import org.springframework.util.Assert; -import org.springframework.util.ClassUtils; import org.springframework.util.StringUtils; /** @@ -100,7 +99,8 @@ import org.springframework.util.StringUtils; * @author David Turanski * @author Marius Bogoevici */ -public class RabbitMessageChannelBinder extends AbstractBinder { +public class RabbitMessageChannelBinder extends AbstractBinder { public static final AnonymousQueue.Base64UrlNamingStrategy ANONYMOUS_GROUP_NAME_GENERATOR = new AnonymousQueue.Base64UrlNamingStrategy("anonymous."); @@ -125,8 +125,6 @@ public class RabbitMessageChannelBinder extends AbstractBinder doBindConsumer(String name, String group, MessageChannel inputChannel, RabbitConsumerProperties properties) { + public Binding doBindConsumer(String name, String group, MessageChannel inputChannel, + RabbitConsumerProperties properties) { boolean anonymousConsumer = !StringUtils.hasText(group); String baseQueueName = anonymousConsumer ? groupedName(name, ANONYMOUS_GROUP_NAME_GENERATOR.generateName()) : groupedName(name, group); @@ -256,7 +253,8 @@ public class RabbitMessageChannelBinder extends AbstractBinder doRegisterConsumer(final String name, String group, MessageChannel moduleInputChannel, Queue queue, final RabbitConsumerProperties properties) { - DefaultBinding consumerBinding = null; - // TODO https://github.com/spring-cloud/spring-cloud-stream/issues/401 - ClassLoader originalClassloader = Thread.currentThread().getContextClassLoader(); - try { - ClassUtils.overrideThreadContextClassLoader(SimpleMessageListenerContainer.class.getClassLoader()); - SimpleMessageListenerContainer listenerContainer = new SimpleMessageListenerContainer( - this.connectionFactory); - listenerContainer.setAcknowledgeMode(properties.getAcknowledgeMode()); - listenerContainer.setChannelTransacted(properties.isTransacted()); - listenerContainer.setDefaultRequeueRejected(properties.isRequeueRejected()); + DefaultBinding consumerBinding; + SimpleMessageListenerContainer listenerContainer = new SimpleMessageListenerContainer( + this.connectionFactory); + listenerContainer.setAcknowledgeMode(properties.getAcknowledgeMode()); + listenerContainer.setChannelTransacted(properties.isTransacted()); + listenerContainer.setDefaultRequeueRejected(properties.isRequeueRejected()); - int concurrency = properties.getConcurrency(); - concurrency = concurrency > 0 ? concurrency : 1; - listenerContainer.setConcurrentConsumers(concurrency); - int maxConcurrency = properties.getMaxConcurrency(); - if (maxConcurrency > concurrency) { - listenerContainer.setMaxConcurrentConsumers(maxConcurrency); - } + int concurrency = properties.getConcurrency(); + concurrency = concurrency > 0 ? concurrency : 1; + listenerContainer.setConcurrentConsumers(concurrency); + int maxConcurrency = properties.getMaxConcurrency(); + if (maxConcurrency > concurrency) { + listenerContainer.setMaxConcurrentConsumers(maxConcurrency); + } - listenerContainer.setPrefetchCount(properties.getPrefetch()); - listenerContainer.setTxSize(properties.getTxSize()); - listenerContainer.setTaskExecutor(new SimpleAsyncTaskExecutor(queue.getName() + "-")); - listenerContainer.setQueues(queue); - int maxAttempts = properties.getMaxAttempts(); - if (maxAttempts > 1 || properties.isRepublishToDlq()) { - RetryOperationsInterceptor retryInterceptor = RetryInterceptorBuilder.stateless() - .maxAttempts(maxAttempts) - .backOffOptions(properties.getBackOffInitialInterval(), - properties.getBackOffMultiplier(), - properties.getBackOffMaxInterval()) - .recoverer(determineRecoverer(name, properties.getPrefix(), properties.isRepublishToDlq())) - .build(); - listenerContainer.setAdviceChain(new Advice[] { retryInterceptor }); + listenerContainer.setPrefetchCount(properties.getPrefetch()); + listenerContainer.setTxSize(properties.getTxSize()); + listenerContainer.setTaskExecutor(new SimpleAsyncTaskExecutor(queue.getName() + "-")); + listenerContainer.setQueues(queue); + int maxAttempts = properties.getMaxAttempts(); + if (maxAttempts > 1 || properties.isRepublishToDlq()) { + RetryOperationsInterceptor retryInterceptor = RetryInterceptorBuilder.stateless() + .maxAttempts(maxAttempts) + .backOffOptions(properties.getBackOffInitialInterval(), + properties.getBackOffMultiplier(), + properties.getBackOffMaxInterval()) + .recoverer(determineRecoverer(name, properties.getPrefix(), properties.isRepublishToDlq())) + .build(); + listenerContainer.setAdviceChain(new Advice[] { retryInterceptor }); + } + listenerContainer.setAfterReceivePostProcessors(this.decompressingPostProcessor); + listenerContainer.setMessagePropertiesConverter(RabbitMessageChannelBinder.inboundMessagePropertiesConverter); + listenerContainer.afterPropertiesSet(); + AmqpInboundChannelAdapter adapter = new AmqpInboundChannelAdapter(listenerContainer); + adapter.setBeanFactory(this.getBeanFactory()); + DirectChannel bridgeToModuleChannel = new DirectChannel(); + bridgeToModuleChannel.setBeanFactory(this.getBeanFactory()); + bridgeToModuleChannel.setBeanName(name + ".bridge"); + adapter.setOutputChannel(bridgeToModuleChannel); + adapter.setBeanName("inbound." + name); + DefaultAmqpHeaderMapper mapper = new DefaultAmqpHeaderMapper(); + mapper.setRequestHeaderNames(properties.getRequestHeaderPatterns()); + mapper.setReplyHeaderNames(properties.getReplyHeaderPatterns()); + adapter.setHeaderMapper(mapper); + adapter.afterPropertiesSet(); + consumerBinding = new DefaultBinding(name, group, moduleInputChannel, adapter) { + @Override + protected void afterUnbind() { + cleanAutoDeclareContext(properties.getPrefix(), name); } - listenerContainer.setAfterReceivePostProcessors(this.decompressingPostProcessor); - listenerContainer.setMessagePropertiesConverter(RabbitMessageChannelBinder.inboundMessagePropertiesConverter); - listenerContainer.afterPropertiesSet(); - AmqpInboundChannelAdapter adapter = new AmqpInboundChannelAdapter(listenerContainer); - adapter.setBeanFactory(this.getBeanFactory()); - DirectChannel bridgeToModuleChannel = new DirectChannel(); - bridgeToModuleChannel.setBeanFactory(this.getBeanFactory()); - bridgeToModuleChannel.setBeanName(name + ".bridge"); - adapter.setOutputChannel(bridgeToModuleChannel); - adapter.setBeanName("inbound." + name); - DefaultAmqpHeaderMapper mapper = new DefaultAmqpHeaderMapper(); - mapper.setRequestHeaderNames(properties.getRequestHeaderPatterns()); - mapper.setReplyHeaderNames(properties.getReplyHeaderPatterns()); - adapter.setHeaderMapper(mapper); - adapter.afterPropertiesSet(); - consumerBinding = new DefaultBinding(name, group, moduleInputChannel, adapter) { - @Override - protected void afterUnbind() { - cleanAutoDeclareContext(properties.getPrefix(), name); - } - }; - ReceivingHandler convertingBridge = new ReceivingHandler(); - convertingBridge.setOutputChannel(moduleInputChannel); - convertingBridge.setBeanName(name + ".convert.bridge"); - convertingBridge.afterPropertiesSet(); - bridgeToModuleChannel.subscribe(convertingBridge); - adapter.start(); - } - finally { - Thread.currentThread().setContextClassLoader(originalClassloader); - } + }; + ReceivingHandler convertingBridge = new ReceivingHandler(); + convertingBridge.setOutputChannel(moduleInputChannel); + convertingBridge.setBeanName(name + ".convert.bridge"); + convertingBridge.afterPropertiesSet(); + bridgeToModuleChannel.subscribe(convertingBridge); + adapter.start(); return consumerBinding; } @@ -426,11 +416,12 @@ public class RabbitMessageChannelBinder extends AbstractBinder