Change default fallback behaviour of MLC if caching connections
This commit is contained in:
@@ -46,7 +46,8 @@ import com.rabbitmq.client.Channel;
|
||||
* @author Mark Fisher
|
||||
* @author Dave Syer
|
||||
*/
|
||||
public class SimpleMessageListenerContainer extends AbstractMessageListenerContainer {
|
||||
public class SimpleMessageListenerContainer extends
|
||||
AbstractMessageListenerContainer {
|
||||
|
||||
public static final long DEFAULT_RECEIVE_TIMEOUT = 1000;
|
||||
|
||||
@@ -77,14 +78,17 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
private CountDownLatch cancellationLock;
|
||||
|
||||
public static interface ContainerDelegate {
|
||||
boolean receiveAndExecute(BlockingQueueConsumer consumer) throws Throwable;
|
||||
boolean receiveAndExecute(BlockingQueueConsumer consumer)
|
||||
throws Throwable;
|
||||
}
|
||||
|
||||
private Advice[] advices = new Advice[0];
|
||||
|
||||
private ContainerDelegate delegate = new ContainerDelegate() {
|
||||
public boolean receiveAndExecute(BlockingQueueConsumer consumer) throws Throwable {
|
||||
return SimpleMessageListenerContainer.this.receiveAndExecute(consumer);
|
||||
public boolean receiveAndExecute(BlockingQueueConsumer consumer)
|
||||
throws Throwable {
|
||||
return SimpleMessageListenerContainer.this
|
||||
.receiveAndExecute(consumer);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -92,16 +96,20 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
|
||||
/**
|
||||
* <p>
|
||||
* Public setter for the {@link Advice} to apply to listener executions. If {@link #setTxSize(int) txSize>1} then
|
||||
* multiple listener executions will all be wrapped in the same advice up to that limit.
|
||||
* Public setter for the {@link Advice} to apply to listener executions. If
|
||||
* {@link #setTxSize(int) txSize>1} then multiple listener executions will
|
||||
* all be wrapped in the same advice up to that limit.
|
||||
* </p>
|
||||
* <p>
|
||||
* If a {@link #setTransactionManager(PlatformTransactionManager) transactionManager} is provided as well, then
|
||||
* separate advice is created for the transaction and applied first in the chain. In that case the advice chain
|
||||
* provided here should not contain a transaction interceptor (otherwise two transactions would be be applied).
|
||||
* If a {@link #setTransactionManager(PlatformTransactionManager)
|
||||
* transactionManager} is provided as well, then separate advice is created
|
||||
* for the transaction and applied first in the chain. In that case the
|
||||
* advice chain provided here should not contain a transaction interceptor
|
||||
* (otherwise two transactions would be be applied).
|
||||
* </p>
|
||||
*
|
||||
* @param advices the advice chain to set
|
||||
* @param advices
|
||||
* the advice chain to set
|
||||
*/
|
||||
public void setAdviceChain(Advice[] advices) {
|
||||
this.advices = advices;
|
||||
@@ -117,12 +125,14 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
/**
|
||||
* Specify the number of concurrent consumers to create. Default is 1.
|
||||
* <p>
|
||||
* Raising the number of concurrent consumers is recommended in order to scale the consumption of messages coming in
|
||||
* from a queue. However, note that any ordering guarantees are lost once multiple consumers are registered. In
|
||||
* general, stick with 1 consumer for low-volume queues.
|
||||
* Raising the number of concurrent consumers is recommended in order to
|
||||
* scale the consumption of messages coming in from a queue. However, note
|
||||
* that any ordering guarantees are lost once multiple consumers are
|
||||
* registered. In general, stick with 1 consumer for low-volume queues.
|
||||
*/
|
||||
public void setConcurrentConsumers(int concurrentConsumers) {
|
||||
Assert.isTrue(concurrentConsumers > 0, "'concurrentConsumers' value must be at least 1 (one)");
|
||||
Assert.isTrue(concurrentConsumers > 0,
|
||||
"'concurrentConsumers' value must be at least 1 (one)");
|
||||
this.concurrentConsumers = concurrentConsumers;
|
||||
}
|
||||
|
||||
@@ -131,12 +141,15 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
}
|
||||
|
||||
/**
|
||||
* The time to wait for workers in milliseconds after the container is stopped, and before the connection is forced
|
||||
* closed. If any workers are active when the shutdown signal comes they will be allowed to finish processing as
|
||||
* long as they can finish within this timeout. Otherwise the connection is closed and messages remain unacked (if
|
||||
* the channel is transactional). Defaults to 5 seconds.
|
||||
* The time to wait for workers in milliseconds after the container is
|
||||
* stopped, and before the connection is forced closed. If any workers are
|
||||
* active when the shutdown signal comes they will be allowed to finish
|
||||
* processing as long as they can finish within this timeout. Otherwise the
|
||||
* connection is closed and messages remain unacked (if the channel is
|
||||
* transactional). Defaults to 5 seconds.
|
||||
*
|
||||
* @param shutdownTimeout the shutdown timeout to set
|
||||
* @param shutdownTimeout
|
||||
* the shutdown timeout to set
|
||||
*/
|
||||
public void setShutdownTimeout(long shutdownTimeout) {
|
||||
this.shutdownTimeout = shutdownTimeout;
|
||||
@@ -148,39 +161,47 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
}
|
||||
|
||||
/**
|
||||
* Tells the broker how many messages to send to each consumer in a single request. Often this can be set quite high
|
||||
* to improve throughput. It should be greater than or equal to {@link #setTxSize(int) the transaction size}.
|
||||
* Tells the broker how many messages to send to each consumer in a single
|
||||
* request. Often this can be set quite high to improve throughput. It
|
||||
* should be greater than or equal to {@link #setTxSize(int) the transaction
|
||||
* size}.
|
||||
*
|
||||
* @param prefetchCount the prefetch count
|
||||
* @param prefetchCount
|
||||
* the prefetch count
|
||||
*/
|
||||
public void setPrefetchCount(int prefetchCount) {
|
||||
this.prefetchCount = prefetchCount;
|
||||
}
|
||||
|
||||
/**
|
||||
* Tells the container how many messages to process in a single transaction (if the channel is transactional). For
|
||||
* best results it should be less than or equal to {@link #setPrefetchCount(int) the prefetch count}.
|
||||
* Tells the container how many messages to process in a single transaction
|
||||
* (if the channel is transactional). For best results it should be less
|
||||
* than or equal to {@link #setPrefetchCount(int) the prefetch count}.
|
||||
*
|
||||
* @param prefetchCount the prefetch count
|
||||
* @param prefetchCount
|
||||
* the prefetch count
|
||||
*/
|
||||
public void setTxSize(int txSize) {
|
||||
this.txSize = txSize;
|
||||
}
|
||||
|
||||
public void setTransactionManager(PlatformTransactionManager transactionManager) {
|
||||
public void setTransactionManager(
|
||||
PlatformTransactionManager transactionManager) {
|
||||
this.transactionManager = transactionManager;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param transactionAttribute the transaction attribute to set
|
||||
* @param transactionAttribute
|
||||
* the transaction attribute to set
|
||||
*/
|
||||
public void setTransactionAttribute(TransactionAttribute transactionAttribute) {
|
||||
public void setTransactionAttribute(
|
||||
TransactionAttribute transactionAttribute) {
|
||||
this.transactionAttribute = transactionAttribute;
|
||||
}
|
||||
|
||||
/**
|
||||
* Avoid the possibility of not configuring the CachingConnectionFactory in sync with the number of concurrent
|
||||
* consumers.
|
||||
* Avoid the possibility of not configuring the CachingConnectionFactory in
|
||||
* sync with the number of concurrent consumers.
|
||||
*/
|
||||
@Override
|
||||
protected void validateConfiguration() {
|
||||
@@ -194,10 +215,6 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
|
||||
if (this.getConnectionFactory() instanceof CachingConnectionFactory) {
|
||||
CachingConnectionFactory cf = (CachingConnectionFactory) getConnectionFactory();
|
||||
if (cf.getChannelCacheSize() < this.concurrentConsumers) {
|
||||
throw new IllegalStateException(
|
||||
"CachingConnectionFactory's channelCacheSize can not be less than the number of concurrentConsumers");
|
||||
}
|
||||
// Default setting
|
||||
if (concurrentConsumers < 1) {
|
||||
concurrentConsumers = 1;
|
||||
@@ -207,8 +224,15 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
logger.info("Setting number of concurrent consumers to CachingConnectionFactory's ChannelCacheSize ["
|
||||
+ cf.getChannelCacheSize() + "]");
|
||||
this.concurrentConsumers = cf.getChannelCacheSize();
|
||||
} else {
|
||||
cf.setChannelCacheSize(1);
|
||||
}
|
||||
}
|
||||
if (cf.getChannelCacheSize() < this.concurrentConsumers) {
|
||||
cf.setChannelCacheSize(this.concurrentConsumers);
|
||||
logger.warn("CachingConnectionFactory's channelCacheSize can not be less than the number of concurrentConsumers so it was reset to match: "
|
||||
+ this.concurrentConsumers);
|
||||
}
|
||||
}
|
||||
if (concurrentConsumers < 1) {
|
||||
concurrentConsumers = 1;
|
||||
@@ -223,8 +247,10 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
if (transactionManager != null) {
|
||||
MatchAlwaysTransactionAttributeSource txAttributeSource = new MatchAlwaysTransactionAttributeSource();
|
||||
txAttributeSource.setTransactionAttribute(transactionAttribute);
|
||||
Advice txAdvice = new TransactionInterceptor(transactionManager, txAttributeSource);
|
||||
factory.addAdvisor(new DefaultPointcutAdvisor(Pointcut.TRUE, txAdvice));
|
||||
Advice txAdvice = new TransactionInterceptor(transactionManager,
|
||||
txAttributeSource);
|
||||
factory.addAdvisor(new DefaultPointcutAdvisor(Pointcut.TRUE,
|
||||
txAdvice));
|
||||
}
|
||||
for (Advice advice : advices) {
|
||||
factory.addAdvisor(new DefaultPointcutAdvisor(Pointcut.TRUE, advice));
|
||||
@@ -247,8 +273,8 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates the specified number of concurrent consumers, in the form of a Rabbit Channel plus associated
|
||||
* MessageConsumer.
|
||||
* Creates the specified number of concurrent consumers, in the form of a
|
||||
* Rabbit Channel plus associated MessageConsumer.
|
||||
*
|
||||
* @throws Exception
|
||||
*/
|
||||
@@ -261,8 +287,9 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
}
|
||||
|
||||
/**
|
||||
* Re-initializes this container's Rabbit message consumers, if not initialized already. Then submits each consumer
|
||||
* to this container's task executor.
|
||||
* Re-initializes this container's Rabbit message consumers, if not
|
||||
* initialized already. Then submits each consumer to this container's task
|
||||
* executor.
|
||||
*
|
||||
* @throws Exception
|
||||
*/
|
||||
@@ -277,7 +304,8 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
}
|
||||
cancellationLock = new CountDownLatch(this.consumers.size());
|
||||
for (BlockingQueueConsumer consumer : this.consumers) {
|
||||
this.taskExecutor.execute(new AsyncMessageProcessingConsumer(consumer, cancellationLock));
|
||||
this.taskExecutor.execute(new AsyncMessageProcessingConsumer(
|
||||
consumer, cancellationLock));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -296,7 +324,8 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
|
||||
try {
|
||||
logger.debug("Waiting for workers to finish.");
|
||||
boolean finished = cancellationLock.await(shutdownTimeout, TimeUnit.MILLISECONDS);
|
||||
boolean finished = cancellationLock.await(shutdownTimeout,
|
||||
TimeUnit.MILLISECONDS);
|
||||
if (finished) {
|
||||
logger.info("Successfully waited for workers to finish.");
|
||||
} else {
|
||||
@@ -316,9 +345,11 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
protected void initializeConsumers() throws IOException {
|
||||
synchronized (this.consumersMonitor) {
|
||||
if (this.consumers == null) {
|
||||
this.consumers = new HashSet<BlockingQueueConsumer>(this.concurrentConsumers);
|
||||
this.consumers = new HashSet<BlockingQueueConsumer>(
|
||||
this.concurrentConsumers);
|
||||
for (int i = 0; i < this.concurrentConsumers; i++) {
|
||||
Channel channel = getTransactionalResourceHolder().getChannel();
|
||||
Channel channel = getTransactionalResourceHolder()
|
||||
.getChannel();
|
||||
BlockingQueueConsumer consumer = createBlockingQueueConsumer(channel);
|
||||
this.consumers.add(consumer);
|
||||
}
|
||||
@@ -327,15 +358,18 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
}
|
||||
|
||||
protected boolean isChannelLocallyTransacted(Channel channel) {
|
||||
return super.isChannelLocallyTransacted(channel) && this.transactionManager == null;
|
||||
return super.isChannelLocallyTransacted(channel)
|
||||
&& this.transactionManager == null;
|
||||
}
|
||||
|
||||
protected BlockingQueueConsumer createBlockingQueueConsumer(final Channel channel) {
|
||||
protected BlockingQueueConsumer createBlockingQueueConsumer(
|
||||
final Channel channel) {
|
||||
BlockingQueueConsumer consumer;
|
||||
String queueNames = getRequiredQueueName();
|
||||
String[] queues = StringUtils.commaDelimitedListToStringArray(queueNames);
|
||||
consumer = new BlockingQueueConsumer(channel, getAcknowledgeMode(), isChannelTransacted(), prefetchCount,
|
||||
queues);
|
||||
String[] queues = StringUtils
|
||||
.commaDelimitedListToStringArray(queueNames);
|
||||
consumer = new BlockingQueueConsumer(channel, getAcknowledgeMode(),
|
||||
isChannelTransacted(), prefetchCount, queues);
|
||||
return consumer;
|
||||
}
|
||||
|
||||
@@ -346,23 +380,29 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
// Need to recycle the channel in this consumer
|
||||
consumer.stop();
|
||||
this.consumers.remove(consumer);
|
||||
Channel channel = getTransactionalResourceHolder().getChannel();
|
||||
Channel channel = getTransactionalResourceHolder()
|
||||
.getChannel();
|
||||
consumer = createBlockingQueueConsumer(channel);
|
||||
this.consumers.add(consumer);
|
||||
} catch (RuntimeException e) {
|
||||
// Ensure consumer counts are correct (another is not going to start because of the exception, but
|
||||
// Ensure consumer counts are correct (another is not going
|
||||
// to start because of the exception, but
|
||||
// we haven't counted down yet)
|
||||
logger.warn("Consumer died on restart. " + e.getClass() + ": " + e.getMessage());
|
||||
logger.warn("Consumer died on restart. " + e.getClass()
|
||||
+ ": " + e.getMessage());
|
||||
cancellationLock.countDown();
|
||||
// Thrown into the void (probably) in a background thread. Oh well, here goes...
|
||||
// Thrown into the void (probably) in a background thread.
|
||||
// Oh well, here goes...
|
||||
throw e;
|
||||
}
|
||||
this.taskExecutor.execute(new AsyncMessageProcessingConsumer(consumer, cancellationLock));
|
||||
this.taskExecutor.execute(new AsyncMessageProcessingConsumer(
|
||||
consumer, cancellationLock));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private boolean receiveAndExecute(BlockingQueueConsumer consumer) throws Throwable {
|
||||
private boolean receiveAndExecute(BlockingQueueConsumer consumer)
|
||||
throws Throwable {
|
||||
|
||||
Channel channel = consumer.getChannel();
|
||||
|
||||
@@ -370,8 +410,8 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
|
||||
ConnectionFactory connectionFactory = getConnectionFactory();
|
||||
if (getAcknowledgeMode().isTransactionAllowed()) {
|
||||
ConnectionFactoryUtils
|
||||
.bindResourceToTransaction(new RabbitResourceHolder(channel), connectionFactory, true);
|
||||
ConnectionFactoryUtils.bindResourceToTransaction(
|
||||
new RabbitResourceHolder(channel), connectionFactory, true);
|
||||
}
|
||||
|
||||
for (int i = 0; i < txSize; i++) {
|
||||
@@ -396,7 +436,8 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
|
||||
private final CountDownLatch latch;
|
||||
|
||||
public AsyncMessageProcessingConsumer(BlockingQueueConsumer consumer, CountDownLatch latch) {
|
||||
public AsyncMessageProcessingConsumer(BlockingQueueConsumer consumer,
|
||||
CountDownLatch latch) {
|
||||
this.consumer = consumer;
|
||||
this.latch = latch;
|
||||
}
|
||||
@@ -407,12 +448,14 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
|
||||
consumer.start();
|
||||
|
||||
// Always better to stop receiving as soon as possible if transactional
|
||||
// Always better to stop receiving as soon as possible if
|
||||
// transactional
|
||||
boolean continuable = false;
|
||||
while (isActive() || continuable) {
|
||||
try {
|
||||
// Will come back false when the queue is drained
|
||||
continuable = proxy.receiveAndExecute(consumer) && !isChannelTransacted();
|
||||
continuable = proxy.receiveAndExecute(consumer)
|
||||
&& !isChannelTransacted();
|
||||
} catch (ListenerExecutionFailedException ex) {
|
||||
// Continue to process, otherwise re-throw
|
||||
}
|
||||
@@ -422,7 +465,9 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
logger.debug("Consumer thread interrupted, processing stopped.");
|
||||
Thread.currentThread().interrupt();
|
||||
} catch (Throwable t) {
|
||||
logger.debug("Consumer received fatal exception, processing stopped.", t);
|
||||
logger.debug(
|
||||
"Consumer received fatal exception, processing stopped.",
|
||||
t);
|
||||
} finally {
|
||||
if (!isActive()) {
|
||||
logger.debug("Cancelling " + consumer);
|
||||
|
||||
Reference in New Issue
Block a user