avoid deadlock in AMQPUrlReceiver hopefully

This commit is contained in:
Noah Levitt
2014-10-17 15:27:08 -07:00
parent 5909042b84
commit d154cfbf36
@@ -25,6 +25,8 @@ import java.util.HashMap;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.logging.Level;
import java.util.logging.Logger;
@@ -104,6 +106,8 @@ public class AMQPUrlReceiver implements Lifecycle, ApplicationListener<CrawlStat
public boolean isRunning() {
return isRunning;
}
private transient Lock lock = new ReentrantLock();
private class StarterRestarter extends Thread {
public StarterRestarter(String name) {
@@ -113,25 +117,29 @@ public class AMQPUrlReceiver implements Lifecycle, ApplicationListener<CrawlStat
@Override
public void run() {
while (!Thread.interrupted()) {
if (!isRunning) {
synchronized (AMQPUrlReceiver.this) {
try {
Consumer consumer = new UrlConsumer(channel());
channel().queueDeclare(getQueueName(), false, false, true, null);
channel().queueBind(getQueueName(), getExchange(), getQueueName());
channel().basicConsume(getQueueName(), false, consumer);
isRunning = true;
logger.info("started AMQP consumer uri=" + getAmqpUri() + " exchange=" + getExchange() + " queueName=" + getQueueName());
} catch (IOException e) {
logger.log(Level.SEVERE, "problem starting AMQP consumer (will try again after 30 seconds)", e);
try {
lock.lockInterruptibly();
if (!isRunning) {
// start up again
synchronized (AMQPUrlReceiver.this) {
try {
Consumer consumer = new UrlConsumer(channel());
channel().queueDeclare(getQueueName(), false, false, true, null);
channel().queueBind(getQueueName(), getExchange(), getQueueName());
channel().basicConsume(getQueueName(), false, consumer);
isRunning = true;
logger.info("started AMQP consumer uri=" + getAmqpUri() + " exchange=" + getExchange() + " queueName=" + getQueueName());
} catch (IOException e) {
logger.log(Level.SEVERE, "problem starting AMQP consumer (will try again after 30 seconds)", e);
}
}
}
}
try {
Thread.sleep(30000);
} catch (InterruptedException e) {
return;
} finally {
lock.unlock();
}
}
}
@@ -140,70 +148,90 @@ public class AMQPUrlReceiver implements Lifecycle, ApplicationListener<CrawlStat
transient private StarterRestarter starterRestarter;
@Override
synchronized public void start() {
// spawn off a thread to start up the amqp consumer, and try to restart it if it dies
if (!isRunning) {
starterRestarter = new StarterRestarter(AMQPUrlReceiver.class.getSimpleName() + "-starter-restarter");
starterRestarter.start();
public void start() {
lock.lock();
try {
// spawn off a thread to start up the amqp consumer, and try to restart it if it dies
if (!isRunning) {
starterRestarter = new StarterRestarter(AMQPUrlReceiver.class.getSimpleName() + "-starter-restarter");
starterRestarter.start();
}
} finally {
lock.unlock();
}
}
@Override
synchronized public void stop() {
logger.info("shutting down");
if (connection != null && connection.isOpen()) {
try {
connection.close();
} catch (IOException e) {
logger.log(Level.SEVERE, "problem closing AMQP connection", e);
public void stop() {
lock.lock();
try {
logger.info("shutting down");
if (connection != null && connection.isOpen()) {
try {
connection.close();
} catch (IOException e) {
logger.log(Level.SEVERE, "problem closing AMQP connection", e);
}
}
}
if (starterRestarter != null && starterRestarter.isAlive()) {
starterRestarter.interrupt();
try {
starterRestarter.join();
} catch (InterruptedException e) {
if (starterRestarter != null && starterRestarter.isAlive()) {
starterRestarter.interrupt();
try {
starterRestarter.join();
} catch (InterruptedException e) {
}
}
starterRestarter = null;
connection = null;
channel = null;
isRunning = false;
} finally {
lock.unlock();
}
starterRestarter = null;
connection = null;
channel = null;
isRunning = false;
}
transient protected Connection connection = null;
transient protected Channel channel = null;
synchronized protected Connection connection() throws IOException {
if (connection != null && !connection.isOpen()) {
logger.warning("connection is closed, creating a new one");
connection = null;
}
if (connection == null) {
ConnectionFactory factory = new ConnectionFactory();
try {
factory.setUri(getAmqpUri());
} catch (Exception e) {
throw new IOException("problem with AMQP uri " + getAmqpUri(), e);
protected Connection connection() throws IOException {
lock.lock();
try {
if (connection != null && !connection.isOpen()) {
logger.warning("connection is closed, creating a new one");
connection = null;
}
connection = factory.newConnection();
}
return connection;
if (connection == null) {
ConnectionFactory factory = new ConnectionFactory();
try {
factory.setUri(getAmqpUri());
} catch (Exception e) {
throw new IOException("problem with AMQP uri " + getAmqpUri(), e);
}
connection = factory.newConnection();
}
return connection;
} finally {
lock.unlock();
}
}
synchronized protected Channel channel() throws IOException {
if (channel != null && !channel.isOpen()) {
logger.warning("channel is not open, creating a new one");
channel = null;
}
protected Channel channel() throws IOException {
lock.lock();
try {
if (channel != null && !channel.isOpen()) {
logger.warning("channel is not open, creating a new one");
channel = null;
}
if (channel == null) {
channel = connection().createChannel();
}
if (channel == null) {
channel = connection().createChannel();
}
return channel;
return channel;
} finally {
lock.unlock();
}
}
// XXX should we be using QueueingConsumer because of possible blocking in