+ * As this bean implements DisposableBean, a bean factory will automatically invoke this on destruction of its
+ * cached singletons.
+ */
+ public final void destroy() {
+ synchronized (this.connectionMonitor) {
+ if (connection != null) {
+ this.connection.destroy();
+ this.connection = null;
+ }
+ }
+ reset();
}
/**
@@ -193,8 +221,7 @@ public class CachingConnectionFactory extends SingleConnectionFactory implements
this.cachedChannelsTransactional.clear();
}
this.active = true;
- super.reset();
- this.targetConnection = null;
+ this.connection = null;
}
@Override
@@ -333,8 +360,15 @@ public class CachingConnectionFactory extends SingleConnectionFactory implements
}
public void close() {
- target.close();
+ }
+
+ public void destroy() {
+ if (this.target != null) {
+ getConnectionListener().onClose(target);
+ RabbitUtils.closeConnection(this.target);
+ }
reset();
+ this.target = null;
}
public boolean isOpen() {
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtils.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtils.java
index d90621a7..fed6e4fb 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtils.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ConnectionFactoryUtils.java
@@ -24,6 +24,7 @@ import org.springframework.transaction.support.TransactionSynchronizationManager
import org.springframework.util.Assert;
import com.rabbitmq.client.Channel;
+
/**
* Helper class for managing a Spring based Rabbit {@link org.springframework.amqp.rabbit.connection.ConnectionFactory},
* in particular for obtaining transactional Rabbit resources for a given ConnectionFactory.
@@ -81,29 +82,26 @@ public class ConnectionFactoryUtils {
final boolean synchedLocalTransactionAllowed) {
RabbitResourceHolder holder = doGetTransactionalResourceHolder(connectionFactory, new ResourceFactory() {
- public Channel getChannel(RabbitResourceHolder holder) {
- return holder.getChannel();
- }
+ 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 connectionFactory.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;
- }
- });
- if (synchedLocalTransactionAllowed) {
- // holder.declareTransactional();
- }
+ public boolean isSynchedLocalTransactionAllowed() {
+ return synchedLocalTransactionAllowed;
+ }
+ });
return holder;
}
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitResourceHolder.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitResourceHolder.java
index aecbdc17..edc7737f 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitResourceHolder.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitResourceHolder.java
@@ -176,18 +176,6 @@ public class RabbitResourceHolder extends ResourceHolderSupport {
}
}
- /**
- * Call this method once the channel {@link Channel#txSelect()} has been called.
- */
- public void declareTransactional() {
- if (!isSynchronizedWithTransaction() && !transactional) {
- for (Channel channel : this.channels) {
- RabbitUtils.declareTransactional(channel);
- }
- this.transactional = !channels.isEmpty();
- }
- }
-
/**
* @return true if the channels in this holder are transactional
*/
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitUtils.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitUtils.java
index cd5ec5f8..47c808dd 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitUtils.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/RabbitUtils.java
@@ -15,13 +15,7 @@ package org.springframework.amqp.rabbit.connection;
import java.io.IOException;
import java.io.UnsupportedEncodingException;
-import java.math.BigDecimal;
import java.net.ConnectException;
-import java.util.Collections;
-import java.util.Date;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
@@ -30,19 +24,11 @@ import org.springframework.amqp.AmqpException;
import org.springframework.amqp.AmqpIOException;
import org.springframework.amqp.AmqpUnsupportedEncodingException;
import org.springframework.amqp.UncategorizedAmqpException;
-import org.springframework.amqp.core.Address;
-import org.springframework.amqp.core.Message;
-import org.springframework.amqp.core.MessageDeliveryMode;
-import org.springframework.amqp.core.MessageProperties;
import org.springframework.util.Assert;
-import org.springframework.util.CollectionUtils;
import com.rabbitmq.client.AMQP;
-import com.rabbitmq.client.AMQP.BasicProperties;
import com.rabbitmq.client.Channel;
-import com.rabbitmq.client.Envelope;
import com.rabbitmq.client.ShutdownSignalException;
-import com.rabbitmq.client.impl.LongString;
/**
* @author Mark Fisher
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/SingleConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/SingleConnectionFactory.java
index 182d12b3..8d63f5d5 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/SingleConnectionFactory.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/SingleConnectionFactory.java
@@ -13,16 +13,9 @@
package org.springframework.amqp.rabbit.connection;
-import java.io.IOException;
-import java.net.InetAddress;
-import java.net.UnknownHostException;
import java.util.List;
-import org.apache.commons.logging.Log;
-import org.apache.commons.logging.LogFactory;
import org.springframework.amqp.AmqpException;
-import org.springframework.beans.factory.DisposableBean;
-import org.springframework.util.Assert;
import org.springframework.util.StringUtils;
import com.rabbitmq.client.Channel;
@@ -35,11 +28,7 @@ import com.rabbitmq.client.Channel;
* @author Mark Pollack
* @author Dave Syer
*/
-public class SingleConnectionFactory implements ConnectionFactory, DisposableBean {
-
- private final Log logger = LogFactory.getLog(getClass());
-
- private final com.rabbitmq.client.ConnectionFactory rabbitConnectionFactory;
+public class SingleConnectionFactory extends AbstractConnectionFactory {
/** Proxy Connection */
private SharedConnectionProxy connection;
@@ -47,8 +36,6 @@ public class SingleConnectionFactory implements ConnectionFactory, DisposableBea
/** Synchronization monitor for the shared Connection */
private final Object connectionMonitor = new Object();
- private final CompositeConnectionListener connectionListener = new CompositeConnectionListener();
-
/**
* Create a new SingleConnectionFactory initializing the hostname to be the value returned from
* InetAddress.getLocalHost(), or "localhost" if getLocalHost() throws an exception.
@@ -79,12 +66,12 @@ public class SingleConnectionFactory implements ConnectionFactory, DisposableBea
* @param port the port number to connect to
*/
public SingleConnectionFactory(String hostname, int port) {
+ super(new com.rabbitmq.client.ConnectionFactory());
if (!StringUtils.hasText(hostname)) {
hostname = getDefaultHostName();
}
- this.rabbitConnectionFactory = new com.rabbitmq.client.ConnectionFactory();
- this.rabbitConnectionFactory.setHost(hostname);
- this.rabbitConnectionFactory.setPort(port);
+ setHost(hostname);
+ setPort(port);
}
/**
@@ -92,52 +79,19 @@ public class SingleConnectionFactory implements ConnectionFactory, DisposableBea
* @param rabbitConnectionFactory the target ConnectionFactory
*/
public SingleConnectionFactory(com.rabbitmq.client.ConnectionFactory rabbitConnectionFactory) {
- Assert.notNull(rabbitConnectionFactory, "Target ConnectionFactory must not be null");
- this.rabbitConnectionFactory = rabbitConnectionFactory;
- }
-
- public void setUsername(String username) {
- this.rabbitConnectionFactory.setUsername(username);
- }
-
- public void setPassword(String password) {
- this.rabbitConnectionFactory.setPassword(password);
- }
-
- public void setHost(String host) {
- this.rabbitConnectionFactory.setHost(host);
- }
-
- public String getHost() {
- return this.rabbitConnectionFactory.getHost();
- }
-
- public void setVirtualHost(String virtualHost) {
- this.rabbitConnectionFactory.setVirtualHost(virtualHost);
- }
-
- public String getVirtualHost() {
- return rabbitConnectionFactory.getVirtualHost();
- }
-
- public void setPort(int port) {
- this.rabbitConnectionFactory.setPort(port);
- }
-
- public int getPort() {
- return this.rabbitConnectionFactory.getPort();
+ super(rabbitConnectionFactory);
}
public void setConnectionListeners(List extends ConnectionListener> listeners) {
- this.connectionListener.setDelegates(listeners);
+ super.setConnectionListeners(listeners);
// If the connection is already alive we assume that the new listeners want to be notified
if (this.connection != null) {
- this.connectionListener.onCreate(this.connection);
+ this.getConnectionListener().onCreate(this.connection);
}
}
public void addConnectionListener(ConnectionListener listener) {
- this.connectionListener.addDelegate(listener);
+ super.addConnectionListener(listener);
// If the connection is already alive we assume that the new listener wants to be notified
if (this.connection != null) {
listener.onCreate(this.connection);
@@ -150,7 +104,7 @@ public class SingleConnectionFactory implements ConnectionFactory, DisposableBea
Connection target = doCreateConnection();
this.connection = new SharedConnectionProxy(target);
// invoke the listener *after* this.connection is assigned
- connectionListener.onCreate(target);
+ getConnectionListener().onCreate(target);
}
}
return this.connection;
@@ -169,13 +123,6 @@ public class SingleConnectionFactory implements ConnectionFactory, DisposableBea
this.connection = null;
}
}
- reset();
- }
-
- /**
- * Default implementation does nothing. Called on {@link #destroy()}.
- */
- protected void reset() {
}
/**
@@ -189,31 +136,9 @@ public class SingleConnectionFactory implements ConnectionFactory, DisposableBea
return connection;
}
- private Connection createBareConnection() {
- try {
- return new SimpleConnection(this.rabbitConnectionFactory.newConnection());
- } catch (IOException e) {
- throw RabbitUtils.convertRabbitAccessException(e);
- }
- }
-
- private String getDefaultHostName() {
- String temp;
- try {
- InetAddress localMachine = InetAddress.getLocalHost();
- temp = localMachine.getHostName();
- logger.debug("Using hostname [" + temp + "] for hostname.");
- } catch (UnknownHostException e) {
- logger.warn("Could not get host name, using 'localhost' as default value", e);
- temp = "localhost";
- }
- return temp;
- }
-
@Override
public String toString() {
- return "SingleConnectionFactory [host=" + rabbitConnectionFactory.getHost() + ", port="
- + rabbitConnectionFactory.getPort() + "]";
+ return "SingleConnectionFactory [host=" + getHost() + ", port=" + getPort() + "]";
}
/**
@@ -235,7 +160,7 @@ public class SingleConnectionFactory implements ConnectionFactory, DisposableBea
if (!isOpen()) {
logger.debug("Detected closed connection. Opening a new one before creating Channel.");
target = createBareConnection();
- connectionListener.onCreate(target);
+ getConnectionListener().onCreate(target);
}
}
}
@@ -248,7 +173,7 @@ public class SingleConnectionFactory implements ConnectionFactory, DisposableBea
public void destroy() {
if (this.target != null) {
- connectionListener.onClose(target);
+ getConnectionListener().onClose(target);
RabbitUtils.closeConnection(this.target);
}
this.target = null;
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java
index bcc80926..a15273eb 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java
@@ -18,7 +18,8 @@ package org.springframework.amqp.rabbit.listener;
import java.io.IOException;
import java.util.ArrayList;
-import java.util.List;
+import java.util.LinkedHashSet;
+import java.util.Set;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
@@ -80,7 +81,7 @@ public class BlockingQueueConsumer {
private final ActiveObjectCounter