From d9bb941af095f4e99cf868586c4d1713ca982543 Mon Sep 17 00:00:00 2001 From: Simon Pamies Date: Wed, 19 Jun 2013 16:08:54 +0200 Subject: [PATCH] Handle reconnection after connection drops for replication listener --- .../test/ConnectionResetAndReconnectTest.java | 7 +-- .../impl/ReplicationMasterDescriptorImpl.java | 2 +- .../server/replication/impl/Replicator.java | 61 +++++++++++++++---- 3 files changed, 53 insertions(+), 17 deletions(-) diff --git a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/ConnectionResetAndReconnectTest.java b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/ConnectionResetAndReconnectTest.java index e533bb3ab43..5f4efbf823a 100644 --- a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/ConnectionResetAndReconnectTest.java +++ b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/ConnectionResetAndReconnectTest.java @@ -14,7 +14,6 @@ import java.util.logging.Logger; import org.junit.Before; import org.junit.Test; -import com.rabbitmq.client.AlreadyClosedException; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; @@ -41,7 +40,7 @@ public class ConnectionResetAndReconnectTest extends AbstractServerReplicationTe @Override public Delivery nextDelivery() throws ShutdownSignalException, ConsumerCancelledException, InterruptedException { if (forceStopDelivery) { - throw new AlreadyClosedException("Could not connect", this); + throw new ShutdownSignalException(false, false, null, null); } return super.nextDelivery(); } @@ -105,10 +104,10 @@ public class ConnectionResetAndReconnectTest extends AbstractServerReplicationTe stopMessagingExchange(); replicaReplicationDescriptor.startToReplicateFrom(masterReplicationDescriptor); Event event = addEventOnMaster(); - Thread.sleep(1000); + Thread.sleep(1000); // wait for master queue to get filled assertNull(replica.getEvent(event.getId())); startMessagingExchange(); - Thread.sleep(2000); // wait for connection to recover + Thread.sleep(3000); // wait for connection to recover assertNotNull(replica.getEvent(event.getId())); } diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationMasterDescriptorImpl.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationMasterDescriptorImpl.java index 3ea274e2594..71a33dd32e3 100755 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationMasterDescriptorImpl.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationMasterDescriptorImpl.java @@ -48,7 +48,7 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip connectionFactory.setPort(port); } Connection connection = connectionFactory.newConnection(); - Channel channel = connection.createChannel(); + Channel channel = connection.createChannel(0); channel.exchangeDeclare(exchangeName, "fanout"); QueueingConsumer consumer = new QueueingConsumer(channel); String queueName = channel.queueDeclare().getQueue(); diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Replicator.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Replicator.java index 7b7d1b03935..5af96a3c75d 100755 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Replicator.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Replicator.java @@ -1,6 +1,7 @@ package com.sap.sailing.server.replication.impl; import java.io.ByteArrayInputStream; +import java.io.IOException; import java.io.ObjectInputStream; import java.util.ArrayList; import java.util.Iterator; @@ -8,7 +9,6 @@ import java.util.List; import java.util.logging.Level; import java.util.logging.Logger; -import com.rabbitmq.client.AlreadyClosedException; import com.rabbitmq.client.QueueingConsumer; import com.rabbitmq.client.QueueingConsumer.Delivery; import com.rabbitmq.client.ShutdownSignalException; @@ -33,10 +33,19 @@ import com.sap.sailing.server.replication.ReplicationMasterDescriptor; public class Replicator implements Runnable { private final static Logger logger = Logger.getLogger(Replicator.class.getName()); + private static final long CHECK_INTERVAL = 2000; // how long (milliseconds) to pause before checking connection again + private static final int CHECK_COUNT = 150; // how long to check, value is CHECK_INTERVAL second steps + private final ReplicationMasterDescriptor master; private final HasRacingEventService racingEventServiceTracker; private final List> queue; - private final QueueingConsumer consumer; + + private QueueingConsumer consumer; + + /** + * How many checks have been performed due to a failing connection? + */ + private int checksPerformed = 0; /** * If the replicator is suspended, messages received are queued. @@ -72,10 +81,13 @@ public class Replicator implements Runnable { @Override public void run() { ClassLoader oldClassLoader = Thread.currentThread().getContextClassLoader(); + while (true) { + boolean noClassLoaderReset = false; try { Delivery delivery = consumer.nextDelivery(); byte[] bytesFromMessage = delivery.getBody(); + checksPerformed = 0; // Set this object's class's class loader as context for de-serialization so that all exported classes // of all required bundles/packages can be deserialized at least Thread.currentThread().setContextClassLoader(getClass().getClassLoader()); @@ -83,21 +95,46 @@ public class Replicator implements Runnable { new ByteArrayInputStream(bytesFromMessage)); RacingEventServiceOperation operation = (RacingEventServiceOperation) ois.readObject(); applyOrQueue(operation); - } catch (AlreadyClosedException ace) { - // ignore this exception and try again after some time - try { - Thread.sleep(1000); - } catch (InterruptedException e) { - e.printStackTrace(); - } } catch (ShutdownSignalException sse) { - logger.info("Received "+sse.getMessage()+". Terminating "+this); - break; + if (sse.isInitiatedByApplication()) { + logger.severe("Application shut down messaging queue for " + this.toString()); + break; + } + + logger.info(sse.getMessage()); + if (checksPerformed <= CHECK_COUNT) { + try { + logger.info("Replication reciever is sleeping because of " + sse.getLocalizedMessage()); + Thread.sleep(CHECK_INTERVAL); + + if (!this.consumer.getChannel().isOpen()) { + /* for a reconnection we need to instantiate a new consumer */ + try { + this.consumer = master.getConsumer(); + Thread.sleep(CHECK_INTERVAL); + checksPerformed += 1; + } catch (IOException eio) { + eio.printStackTrace(); + } + } + } catch (InterruptedException eir) { + eir.printStackTrace(); + } + checksPerformed += 1; + noClassLoaderReset = true; + continue; + } else { + logger.severe("Grace time (" + CHECK_COUNT*(CHECK_INTERVAL/1000) + "secs) is over. Terminating replication listener " + this.toString()); + // XXX: Also make sure that all handlers get notifications about this + break; + } } catch (Exception e) { logger.info("Exception while processing replica: "+e.getMessage()); logger.log(Level.SEVERE, "run", e); } finally { - Thread.currentThread().setContextClassLoader(oldClassLoader); + if (noClassLoaderReset == false) { + Thread.currentThread().setContextClassLoader(oldClassLoader); + } } } }