Revamped functionality of removing a queue. It works now correctly and

removes queues.
This commit is contained in:
Simon Pamies committed 2013-06-21 01:50:58 +02:00
1 parent 464d7d2051
commit 25a46ef312
4 files changed
+24 -7

No files matched your search

@@ -36,4 +36,6 @@ public interface ReplicationMasterDescriptor {
QueueingConsumer getConsumer() throws IOException;
String getExchangeName();
void stopConnection();
}
@@ -23,6 +23,8 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
private final int messagingPort;
private final String queueName;
private QueueingConsumer consumer;
/**
* @param messagingPort 0 means use default port
*/
@@ -32,6 +34,7 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
this.messagingPort = messagingPort;
this.exchangeName = exchangeName;
this.queueName = queueName;
this.consumer = null;
}
@Override
@@ -56,7 +59,7 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
}
@Override
public QueueingConsumer getConsumer() throws IOException {
public synchronized QueueingConsumer getConsumer() throws IOException {
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost(getHostname());
int port = getMessagingPort();
@@ -108,8 +111,24 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
// from now on we get all new messages that the exchange is getting from producer
channel.queueBind(queueName, exchangeName, "");
channel.basicConsume(queueName, /* auto-ack */ true, consumer);
this.consumer = consumer;
return consumer;
}
@Override
public synchronized void stopConnection() {
try {
if (consumer != null) {
// make sure to remove queue in order to avoid any exchanges filling it with messages
consumer.getChannel().queueUnbind(queueName, exchangeName, "");
consumer.getChannel().queueDelete(queueName);
consumer.getChannel().getConnection().close(1);
}
} catch (Exception ex) {
// ignore any exception during abort. close can yield a broad
// number of exceptions that we don't want to know or to log.
}
}
/**
* @return 0 means use default port
@@ -300,7 +300,6 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
if (replicator != null) {
replicator.stop();
deregisterReplicaWithMaster(descriptor);
descriptor.getConsumer().getChannel().close();
replicatingFromMaster = null;
replicaUUIDs.clear();
@@ -149,6 +149,7 @@ public class Replicator implements Runnable {
Thread.currentThread().setContextClassLoader(oldClassLoader);
}
}
logger.info("Stopped replicator thread. This server will no longer receive events from a master.");
}
@@ -210,11 +211,7 @@ public class Replicator implements Runnable {
}
logger.info("Signaled Replicator thread to stop asap.");
stopped = true;
try {
master.getConsumer().getChannel().getConnection().close(1);
} catch (Exception ex) {
// ignore any exception during abort.
}
master.stopConnection();
}
public synchronized boolean isBeingStopped() {