From 10e0f262b31a8110c0d009a99f0d840c1003941e Mon Sep 17 00:00:00 2001 From: Noah Levitt Date: Thu, 24 Apr 2014 20:31:40 -0700 Subject: [PATCH 1/2] ARI-3765 - make AMQPPublishProcessor gracefully handle amqp server going up and down --- .../archive/modules/AMQPPublishProcessor.java | 19 +++++++++++++++++-- 1 file changed, 17 insertions(+), 2 deletions(-) diff --git a/contrib/src/main/java/org/archive/modules/AMQPPublishProcessor.java b/contrib/src/main/java/org/archive/modules/AMQPPublishProcessor.java index 02935cb8..3646f399 100644 --- a/contrib/src/main/java/org/archive/modules/AMQPPublishProcessor.java +++ b/contrib/src/main/java/org/archive/modules/AMQPPublishProcessor.java @@ -24,6 +24,7 @@ import static org.archive.modules.CoreAttributeConstants.A_HERITABLE_KEYS; import java.io.IOException; import java.util.HashMap; import java.util.Set; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.logging.Level; import java.util.logging.Logger; @@ -176,8 +177,12 @@ public class AMQPPublishProcessor extends Processor { } protected synchronized Channel channel() { + if (threadChannel.get() != null && !threadChannel.get().isOpen()) { + threadChannel.set(null); + } + if (threadChannel.get() == null) { - if (connection == null) { + if (connection == null || !connection.isOpen()) { connect(); } try { @@ -190,14 +195,24 @@ public class AMQPPublishProcessor extends Processor { } return threadChannel.get(); } + + private AtomicBoolean serverLooksDown = new AtomicBoolean(false); private synchronized void connect() { ConnectionFactory factory = new ConnectionFactory(); try { factory.setUri(getAmqpUri()); connection = factory.newConnection(); + boolean wasDown = serverLooksDown.getAndSet(false); + if (wasDown) { + logger.info(getAmqpUri() + " is back up, connected successfully!"); + } } catch (Exception e) { - logger.log(Level.SEVERE, "Attempting to connect to AMQP server failed!", e); + connection = null; + boolean wasAlreadyDown = serverLooksDown.getAndSet(true); + if (!wasAlreadyDown) { + logger.log(Level.SEVERE, "Attempting to connect to AMQP server failed!", e); + } } } From 381442dfb60ac2157795177555529c97d4b84136 Mon Sep 17 00:00:00 2001 From: Noah Levitt Date: Thu, 24 Apr 2014 21:23:03 -0700 Subject: [PATCH 2/2] ARI-3765 - make AMQPUrlReceiver gracefully handle amqp server going up and down --- .../crawler/frontier/AMQPUrlReceiver.java | 72 ++++++++++++------- 1 file changed, 47 insertions(+), 25 deletions(-) diff --git a/contrib/src/main/java/org/archive/crawler/frontier/AMQPUrlReceiver.java b/contrib/src/main/java/org/archive/crawler/frontier/AMQPUrlReceiver.java index c1c37f07..4074a842 100644 --- a/contrib/src/main/java/org/archive/crawler/frontier/AMQPUrlReceiver.java +++ b/contrib/src/main/java/org/archive/crawler/frontier/AMQPUrlReceiver.java @@ -51,6 +51,7 @@ import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.Consumer; import com.rabbitmq.client.DefaultConsumer; import com.rabbitmq.client.Envelope; +import com.rabbitmq.client.ShutdownSignalException; /** * @contributor nlevitt @@ -104,32 +105,43 @@ public class AMQPUrlReceiver implements Lifecycle, ApplicationListener