mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-10-11 06:40:53 +00:00
Handle reconnection after connection drops for replication listener
This commit is contained in:
1 parent
d337e7206e
commit
d9bb941af0
3 files changed
+53
-17
No files matched your search
+1
-1
@@ -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();
|
||||
|
||||
+49
-12
@@ -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<RacingEventServiceOperation<?>> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user