AMQP-113, AMQP-112: backoff and sleep when listener container fails to connect
This commit is contained in:
@@ -38,6 +38,12 @@
|
||||
<groupId>org.springframework</groupId>
|
||||
<artifactId>spring-test</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework</groupId>
|
||||
<artifactId>spring-retry</artifactId>
|
||||
<version>1.0.0.RC1</version>
|
||||
<optional>true</optional>
|
||||
</dependency>
|
||||
|
||||
<!-- Other -->
|
||||
<dependency>
|
||||
|
||||
@@ -13,11 +13,8 @@
|
||||
|
||||
package org.springframework.amqp.rabbit.config;
|
||||
|
||||
import org.w3c.dom.Element;
|
||||
import org.w3c.dom.Node;
|
||||
import org.w3c.dom.NodeList;
|
||||
|
||||
import org.springframework.amqp.core.AcknowledgeMode;
|
||||
import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer;
|
||||
import org.springframework.beans.factory.config.BeanDefinition;
|
||||
import org.springframework.beans.factory.config.RuntimeBeanReference;
|
||||
import org.springframework.beans.factory.parsing.BeanComponentDefinition;
|
||||
@@ -26,6 +23,9 @@ import org.springframework.beans.factory.support.RootBeanDefinition;
|
||||
import org.springframework.beans.factory.xml.BeanDefinitionParser;
|
||||
import org.springframework.beans.factory.xml.ParserContext;
|
||||
import org.springframework.util.StringUtils;
|
||||
import org.w3c.dom.Element;
|
||||
import org.w3c.dom.Node;
|
||||
import org.w3c.dom.NodeList;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
@@ -162,8 +162,7 @@ class ListenerContainerParser implements BeanDefinitionParser {
|
||||
}
|
||||
|
||||
private BeanDefinition parseContainer(Element listenerEle, Element containerEle, ParserContext parserContext) {
|
||||
RootBeanDefinition containerDef = new RootBeanDefinition(
|
||||
"org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer");
|
||||
RootBeanDefinition containerDef = new RootBeanDefinition(SimpleMessageListenerContainer.class);
|
||||
containerDef.setSource(parserContext.extractSource(containerEle));
|
||||
|
||||
String connectionFactoryBeanName = "rabbitConnectionFactory";
|
||||
|
||||
@@ -71,37 +71,41 @@ public class ConnectionFactoryUtils {
|
||||
|
||||
/**
|
||||
* Obtain a RabbitMQ Channel that is synchronized with the current transaction, if any.
|
||||
* @param cf the ConnectionFactory to obtain a Channel for
|
||||
* @param connectionFactory the ConnectionFactory to obtain a Channel for
|
||||
* @param synchedLocalTransactionAllowed whether to allow for a local RabbitMQ transaction that is synchronized with
|
||||
* a Spring-managed transaction (where the main transaction might be a JDBC-based one for a specific DataSource, for
|
||||
* example), with the RabbitMQ transaction committing right after the main transaction. If not allowed, the given
|
||||
* ConnectionFactory needs to handle transaction enlistment underneath the covers.
|
||||
* @return the transactional Channel, or <code>null</code> if none found
|
||||
*/
|
||||
public static RabbitResourceHolder getTransactionalResourceHolder(final ConnectionFactory cf,
|
||||
public static RabbitResourceHolder getTransactionalResourceHolder(final ConnectionFactory connectionFactory,
|
||||
final boolean synchedLocalTransactionAllowed) {
|
||||
|
||||
return doGetTransactionalResourceHolder(cf, new ResourceFactory() {
|
||||
public Channel getChannel(RabbitResourceHolder holder) {
|
||||
return holder.getChannel();
|
||||
}
|
||||
RabbitResourceHolder holder = doGetTransactionalResourceHolder(connectionFactory, new ResourceFactory() {
|
||||
public Channel getChannel(RabbitResourceHolder holder) {
|
||||
return holder.getChannel();
|
||||
}
|
||||
|
||||
public Connection getConnection(RabbitResourceHolder holder) {
|
||||
return holder.getConnection();
|
||||
}
|
||||
public Connection getConnection(RabbitResourceHolder holder) {
|
||||
return holder.getConnection();
|
||||
}
|
||||
|
||||
public Connection createConnection() throws IOException {
|
||||
return cf.createConnection();
|
||||
}
|
||||
public Connection createConnection() throws IOException {
|
||||
return connectionFactory.createConnection();
|
||||
}
|
||||
|
||||
public Channel createChannel(Connection con) throws IOException {
|
||||
return con.createChannel(synchedLocalTransactionAllowed);
|
||||
}
|
||||
public Channel createChannel(Connection con) throws IOException {
|
||||
return con.createChannel(synchedLocalTransactionAllowed);
|
||||
}
|
||||
|
||||
public boolean isSynchedLocalTransactionAllowed() {
|
||||
return synchedLocalTransactionAllowed;
|
||||
}
|
||||
});
|
||||
public boolean isSynchedLocalTransactionAllowed() {
|
||||
return synchedLocalTransactionAllowed;
|
||||
}
|
||||
});
|
||||
if (synchedLocalTransactionAllowed) {
|
||||
holder.declareTransactional();
|
||||
}
|
||||
return holder;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -12,6 +12,8 @@ import org.springframework.amqp.AmqpException;
|
||||
import org.springframework.amqp.core.AcknowledgeMode;
|
||||
import org.springframework.amqp.core.Message;
|
||||
import org.springframework.amqp.core.MessageProperties;
|
||||
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
|
||||
import org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils;
|
||||
import org.springframework.amqp.rabbit.support.RabbitUtils;
|
||||
|
||||
import com.rabbitmq.client.AMQP;
|
||||
@@ -23,9 +25,7 @@ import com.rabbitmq.client.ShutdownSignalException;
|
||||
import com.rabbitmq.utility.Utility;
|
||||
|
||||
/**
|
||||
* Variation on QueueingConsumer in RabbitMQ, uses 'put' instead of 'add' and stored a reference to the consumerTag that
|
||||
* was returned when this Consumer was registered with the channel so as to make it easy to close the consumer when
|
||||
* shutting down.
|
||||
* Specialized consumer encapsulating knowledge of the broker connections and having its own lifecycle (start and stop).
|
||||
*
|
||||
* @author Mark Pollack
|
||||
* @author Dave Syer
|
||||
@@ -35,6 +35,7 @@ public class BlockingQueueConsumer {
|
||||
|
||||
private static Log logger = LogFactory.getLog(BlockingQueueConsumer.class);
|
||||
|
||||
// This must be an unbounded queue or we risk blocking the Connection thread.
|
||||
private final BlockingQueue<Delivery> queue = new LinkedBlockingQueue<Delivery>();
|
||||
|
||||
// When this is non-null the connection has been closed (should never happen in normal operation).
|
||||
@@ -46,22 +47,27 @@ public class BlockingQueueConsumer {
|
||||
|
||||
private final boolean transactional;
|
||||
|
||||
private final Channel channel;
|
||||
private Channel channel;
|
||||
|
||||
private InternalConsumer consumer;
|
||||
|
||||
private final AtomicBoolean cancelled = new AtomicBoolean(false);
|
||||
|
||||
private final InternalConsumer consumer;
|
||||
|
||||
private final AcknowledgeMode acknowledgeMode;
|
||||
|
||||
public BlockingQueueConsumer(Channel channel, AcknowledgeMode acknowledgeMode, boolean transactional,
|
||||
int prefetchCount, String... queues) {
|
||||
this.channel = channel;
|
||||
private final ConnectionFactory connectionFactory;
|
||||
|
||||
/**
|
||||
* Create a consumer. The consumer must not attempt to use the connection factory or communicate with the broker
|
||||
* until it is started.
|
||||
*/
|
||||
public BlockingQueueConsumer(ConnectionFactory connectionFactory, AcknowledgeMode acknowledgeMode,
|
||||
boolean transactional, int prefetchCount, String... queues) {
|
||||
this.connectionFactory = connectionFactory;
|
||||
this.acknowledgeMode = acknowledgeMode;
|
||||
this.transactional = transactional;
|
||||
this.prefetchCount = prefetchCount;
|
||||
this.queues = queues;
|
||||
this.consumer = new InternalConsumer(channel);
|
||||
}
|
||||
|
||||
public Channel getChannel() {
|
||||
@@ -136,6 +142,9 @@ public class BlockingQueueConsumer {
|
||||
}
|
||||
|
||||
public void start() throws AmqpException {
|
||||
this.channel = ConnectionFactoryUtils.getTransactionalResourceHolder(connectionFactory, transactional)
|
||||
.getChannel();
|
||||
this.consumer = new InternalConsumer(channel);
|
||||
try {
|
||||
// Set basicQos before calling basicConsume (it is ignored if we are not transactional and the broker will
|
||||
// send blocks of 100 messages)
|
||||
@@ -155,14 +164,8 @@ public class BlockingQueueConsumer {
|
||||
public void stop() {
|
||||
cancelled.set(true);
|
||||
logger.debug("Closing Rabbit Channel: " + channel);
|
||||
try {
|
||||
if (consumer != null && consumer.getChannel() != null && consumer.getConsumerTag() != null) {
|
||||
RabbitUtils.closeMessageConsumer(consumer.getChannel(), consumer.getConsumerTag(), transactional);
|
||||
} catch (AmqpException e) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.info("Could not close message consumer on shutdown", e);
|
||||
} else {
|
||||
logger.info("Could not close message consumer on shutdown (" + e.getClass() + "): " + e.getMessage());
|
||||
}
|
||||
}
|
||||
// This one never throws exceptions...
|
||||
RabbitUtils.closeChannel(channel);
|
||||
@@ -234,8 +237,8 @@ public class BlockingQueueConsumer {
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "Consumer: tag=[" + consumer.getConsumerTag() + "], channel=" + channel + ", acknowledgeMode="
|
||||
+ acknowledgeMode + " local queue size=" + queue.size();
|
||||
return "Consumer: tag=[" + (consumer != null ? consumer.getConsumerTag() : null) + "], channel=" + channel
|
||||
+ ", acknowledgeMode=" + acknowledgeMode + " local queue size=" + queue.size();
|
||||
}
|
||||
|
||||
}
|
||||
@@ -13,7 +13,6 @@
|
||||
|
||||
package org.springframework.amqp.rabbit.listener;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.HashSet;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
@@ -21,6 +20,7 @@ import java.util.concurrent.Executor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import org.aopalliance.aop.Advice;
|
||||
import org.springframework.amqp.AmqpException;
|
||||
import org.springframework.amqp.core.Message;
|
||||
import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
|
||||
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
|
||||
@@ -46,14 +46,18 @@ 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;
|
||||
|
||||
private static final int DEFAULT_PREFETCH_COUNT = 1;
|
||||
public static final int DEFAULT_PREFETCH_COUNT = 1;
|
||||
|
||||
private static final long DEFAULT_SHUTDOWN_TIMEOUT = 5000;
|
||||
public static final long DEFAULT_SHUTDOWN_TIMEOUT = 5000;
|
||||
|
||||
/**
|
||||
* The default recovery interval: 5000 ms = 5 seconds.
|
||||
*/
|
||||
public static final long DEFAULT_RECOVERY_INTERVAL = 5000;
|
||||
|
||||
private volatile int prefetchCount = DEFAULT_PREFETCH_COUNT;
|
||||
|
||||
@@ -67,6 +71,8 @@ public class SimpleMessageListenerContainer extends
|
||||
|
||||
private long shutdownTimeout = DEFAULT_SHUTDOWN_TIMEOUT;
|
||||
|
||||
private long recoveryInterval = DEFAULT_RECOVERY_INTERVAL;
|
||||
|
||||
private Set<BlockingQueueConsumer> consumers;
|
||||
|
||||
private final Object consumersMonitor = new Object();
|
||||
@@ -78,17 +84,14 @@ public class SimpleMessageListenerContainer extends
|
||||
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);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -96,25 +99,30 @@ public class SimpleMessageListenerContainer extends
|
||||
|
||||
/**
|
||||
* <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;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify the interval between recovery attempts, in <b>milliseconds</b>. The default is 5000 ms, that is, 5
|
||||
* seconds.
|
||||
* @see #handleConsumerStartupFailure
|
||||
*/
|
||||
public void setRecoveryInterval(long recoveryInterval) {
|
||||
this.recoveryInterval = recoveryInterval;
|
||||
}
|
||||
|
||||
public SimpleMessageListenerContainer() {
|
||||
}
|
||||
|
||||
@@ -125,14 +133,12 @@ public class SimpleMessageListenerContainer extends
|
||||
/**
|
||||
* 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;
|
||||
}
|
||||
|
||||
@@ -141,15 +147,12 @@ public class SimpleMessageListenerContainer extends
|
||||
}
|
||||
|
||||
/**
|
||||
* 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;
|
||||
@@ -161,47 +164,39 @@ public class SimpleMessageListenerContainer extends
|
||||
}
|
||||
|
||||
/**
|
||||
* 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() {
|
||||
@@ -239,7 +234,7 @@ public class SimpleMessageListenerContainer extends
|
||||
}
|
||||
}
|
||||
|
||||
public void initializeProxy() {
|
||||
private void initializeProxy() {
|
||||
if (advices.length == 0 && transactionManager == null) {
|
||||
return;
|
||||
}
|
||||
@@ -247,12 +242,10 @@ public class SimpleMessageListenerContainer extends
|
||||
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) {
|
||||
for (Advice advice : getAdvices()) {
|
||||
factory.addAdvisor(new DefaultPointcutAdvisor(Pointcut.TRUE, advice));
|
||||
}
|
||||
factory.setProxyTargetClass(false);
|
||||
@@ -273,8 +266,8 @@ public class SimpleMessageListenerContainer extends
|
||||
}
|
||||
|
||||
/**
|
||||
* 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
|
||||
*/
|
||||
@@ -286,16 +279,25 @@ public class SimpleMessageListenerContainer extends
|
||||
return (int) cancellationLock.getCount();
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void establishSharedConnection() throws Exception {
|
||||
try {
|
||||
super.establishSharedConnection();
|
||||
} catch (AmqpException e) {
|
||||
// Try to recover later...
|
||||
logger.info("Could not start message listener container. Consumer threads will attempt to reconnect. Exception ("
|
||||
+ e.getClass().getName() + "): " + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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
|
||||
*/
|
||||
protected void doStart() throws Exception {
|
||||
super.doStart();
|
||||
establishSharedConnection();
|
||||
initializeConsumers();
|
||||
synchronized (this.consumersMonitor) {
|
||||
if (this.consumers == null) {
|
||||
@@ -304,8 +306,7 @@ public class SimpleMessageListenerContainer extends
|
||||
}
|
||||
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));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -324,8 +325,7 @@ public class SimpleMessageListenerContainer extends
|
||||
|
||||
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 {
|
||||
@@ -342,15 +342,12 @@ public class SimpleMessageListenerContainer extends
|
||||
|
||||
}
|
||||
|
||||
protected void initializeConsumers() throws IOException {
|
||||
protected void initializeConsumers() {
|
||||
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();
|
||||
BlockingQueueConsumer consumer = createBlockingQueueConsumer(channel);
|
||||
BlockingQueueConsumer consumer = createBlockingQueueConsumer();
|
||||
this.consumers.add(consumer);
|
||||
}
|
||||
}
|
||||
@@ -358,18 +355,14 @@ public class SimpleMessageListenerContainer extends
|
||||
}
|
||||
|
||||
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() {
|
||||
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(getConnectionFactory(), getAcknowledgeMode(), isChannelTransacted(), prefetchCount, queues);
|
||||
return consumer;
|
||||
}
|
||||
|
||||
@@ -380,29 +373,24 @@ public class SimpleMessageListenerContainer extends
|
||||
// Need to recycle the channel in this consumer
|
||||
consumer.stop();
|
||||
this.consumers.remove(consumer);
|
||||
Channel channel = getTransactionalResourceHolder()
|
||||
.getChannel();
|
||||
consumer = createBlockingQueueConsumer(channel);
|
||||
consumer = createBlockingQueueConsumer();
|
||||
this.consumers.add(consumer);
|
||||
} catch (RuntimeException e) {
|
||||
// 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...
|
||||
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();
|
||||
|
||||
@@ -410,8 +398,8 @@ public class SimpleMessageListenerContainer extends
|
||||
|
||||
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++) {
|
||||
@@ -430,14 +418,17 @@ public class SimpleMessageListenerContainer extends
|
||||
|
||||
}
|
||||
|
||||
protected Advice[] getAdvices() {
|
||||
return advices;
|
||||
}
|
||||
|
||||
private class AsyncMessageProcessingConsumer implements Runnable {
|
||||
|
||||
private final BlockingQueueConsumer consumer;
|
||||
|
||||
private final CountDownLatch latch;
|
||||
|
||||
public AsyncMessageProcessingConsumer(BlockingQueueConsumer consumer,
|
||||
CountDownLatch latch) {
|
||||
public AsyncMessageProcessingConsumer(BlockingQueueConsumer consumer, CountDownLatch latch) {
|
||||
this.consumer = consumer;
|
||||
this.latch = latch;
|
||||
}
|
||||
@@ -446,7 +437,12 @@ public class SimpleMessageListenerContainer extends
|
||||
|
||||
try {
|
||||
|
||||
consumer.start();
|
||||
try {
|
||||
consumer.start();
|
||||
} catch (Throwable t) {
|
||||
handleStartupFailure(t);
|
||||
throw t;
|
||||
}
|
||||
|
||||
// Always better to stop receiving as soon as possible if
|
||||
// transactional
|
||||
@@ -454,8 +450,7 @@ public class SimpleMessageListenerContainer extends
|
||||
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
|
||||
}
|
||||
@@ -465,16 +460,22 @@ public class SimpleMessageListenerContainer extends
|
||||
logger.debug("Consumer thread interrupted, processing stopped.");
|
||||
Thread.currentThread().interrupt();
|
||||
} catch (Throwable t) {
|
||||
logger.debug(
|
||||
"Consumer received fatal exception, processing stopped.",
|
||||
t);
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("Consumer received fatal exception, processing stopped", t);
|
||||
} else {
|
||||
logger.warn("Consumer received fatal exception, processing stopped: " + t);
|
||||
}
|
||||
} finally {
|
||||
if (!isActive()) {
|
||||
logger.debug("Cancelling " + consumer);
|
||||
latch.countDown();
|
||||
consumer.stop();
|
||||
try {
|
||||
consumer.stop();
|
||||
} catch (AmqpException e) {
|
||||
logger.info("Could not stop message consumer on shutdown", e);
|
||||
}
|
||||
} else {
|
||||
logger.debug("Restarting " + consumer);
|
||||
logger.info("Restarting " + consumer);
|
||||
restart(consumer);
|
||||
}
|
||||
}
|
||||
@@ -483,4 +484,20 @@ public class SimpleMessageListenerContainer extends
|
||||
|
||||
}
|
||||
|
||||
/**
|
||||
* Give the container a chance to recover from consumer startup failure, e.g. if teh broker is down.
|
||||
*
|
||||
* @param t the exception that stopped the startup
|
||||
* @throws Exception if the shared connection still can't be established
|
||||
*/
|
||||
protected void handleStartupFailure(Throwable t) throws Exception {
|
||||
try {
|
||||
Thread.sleep(recoveryInterval);
|
||||
establishSharedConnection();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException("Unrecoverable interruption on consumer restart");
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -110,9 +110,6 @@ public abstract class RabbitAccessor implements InitializingBean {
|
||||
|
||||
protected RabbitResourceHolder getTransactionalResourceHolder() {
|
||||
RabbitResourceHolder holder = ConnectionFactoryUtils.getTransactionalResourceHolder(this.connectionFactory, isChannelTransacted());
|
||||
if (isChannelTransacted()) {
|
||||
holder.declareTransactional();
|
||||
}
|
||||
return holder;
|
||||
}
|
||||
|
||||
|
||||
@@ -13,6 +13,7 @@ import org.springframework.amqp.core.Queue;
|
||||
import org.springframework.amqp.core.TopicExchange;
|
||||
import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
|
||||
import org.springframework.amqp.rabbit.listener.BlockingQueueConsumer;
|
||||
import org.springframework.amqp.rabbit.support.RabbitAccessor;
|
||||
import org.springframework.amqp.rabbit.test.BrokerRunning;
|
||||
import org.springframework.amqp.support.converter.SimpleMessageConverter;
|
||||
|
||||
@@ -40,7 +41,7 @@ public class RabbitBindingIntegrationTests {
|
||||
template.execute(new ChannelCallback<Void>() {
|
||||
public Void doInRabbit(Channel channel) throws Exception {
|
||||
|
||||
BlockingQueueConsumer consumer = createConsumer(channel);
|
||||
BlockingQueueConsumer consumer = createConsumer(template);
|
||||
String tag = consumer.getConsumerTag();
|
||||
assertNotNull(tag);
|
||||
|
||||
@@ -79,7 +80,7 @@ public class RabbitBindingIntegrationTests {
|
||||
template.execute(new ChannelCallback<Void>() {
|
||||
public Void doInRabbit(Channel channel) throws Exception {
|
||||
|
||||
BlockingQueueConsumer consumer = createConsumer(channel);
|
||||
BlockingQueueConsumer consumer = createConsumer(template);
|
||||
String tag = consumer.getConsumerTag();
|
||||
assertNotNull(tag);
|
||||
|
||||
@@ -122,7 +123,7 @@ public class RabbitBindingIntegrationTests {
|
||||
BlockingQueueConsumer consumer = template.execute(new ChannelCallback<BlockingQueueConsumer>() {
|
||||
public BlockingQueueConsumer doInRabbit(Channel channel) throws Exception {
|
||||
|
||||
BlockingQueueConsumer consumer = createConsumer(channel);
|
||||
BlockingQueueConsumer consumer = createConsumer(template);
|
||||
String tag = consumer.getConsumerTag();
|
||||
assertNotNull(tag);
|
||||
|
||||
@@ -156,7 +157,7 @@ public class RabbitBindingIntegrationTests {
|
||||
template.execute(new ChannelCallback<Void>() {
|
||||
public Void doInRabbit(Channel channel) throws Exception {
|
||||
|
||||
BlockingQueueConsumer consumer = createConsumer(channel);
|
||||
BlockingQueueConsumer consumer = createConsumer(template);
|
||||
String tag = consumer.getConsumerTag();
|
||||
assertNotNull(tag);
|
||||
|
||||
@@ -176,7 +177,7 @@ public class RabbitBindingIntegrationTests {
|
||||
template.execute(new ChannelCallback<Void>() {
|
||||
public Void doInRabbit(Channel channel) throws Exception {
|
||||
|
||||
BlockingQueueConsumer consumer = createConsumer(channel);
|
||||
BlockingQueueConsumer consumer = createConsumer(template);
|
||||
String tag = consumer.getConsumerTag();
|
||||
assertNotNull(tag);
|
||||
|
||||
@@ -208,7 +209,7 @@ public class RabbitBindingIntegrationTests {
|
||||
template.execute(new ChannelCallback<Void>() {
|
||||
public Void doInRabbit(Channel channel) throws Exception {
|
||||
|
||||
BlockingQueueConsumer consumer = createConsumer(channel);
|
||||
BlockingQueueConsumer consumer = createConsumer(template);
|
||||
String tag = consumer.getConsumerTag();
|
||||
assertNotNull(tag);
|
||||
|
||||
@@ -227,8 +228,8 @@ public class RabbitBindingIntegrationTests {
|
||||
|
||||
}
|
||||
|
||||
private BlockingQueueConsumer createConsumer(Channel channel) {
|
||||
BlockingQueueConsumer consumer = new BlockingQueueConsumer(channel, AcknowledgeMode.AUTO, true, 1, queue.getName());
|
||||
private BlockingQueueConsumer createConsumer(RabbitAccessor accessor) {
|
||||
BlockingQueueConsumer consumer = new BlockingQueueConsumer(accessor.getConnectionFactory(), AcknowledgeMode.AUTO, true, 1, queue.getName());
|
||||
consumer.start();
|
||||
return consumer;
|
||||
}
|
||||
|
||||
@@ -41,7 +41,7 @@
|
||||
|
||||
<rabbit:queue name="bar" />
|
||||
|
||||
<rabbit:admin id="admin-test" rabbit-connection-factory="connectionFactory"/>
|
||||
<rabbit:admin id="admin-test" connection-factory="connectionFactory"/>
|
||||
|
||||
<bean id="connectionFactory" class="org.springframework.amqp.rabbit.connection.CachingConnectionFactory"/>
|
||||
|
||||
|
||||
Reference in New Issue
Block a user