a compiling, yet entirely untested version preparing for RabbitMQ-based instead of ActiveMQ-based replication

This commit is contained in:
Axel Uhl
2012-07-11 15:01:45 +02:00
parent 193b956f09
commit 7d4ac52bac
23 changed files with 236 additions and 502 deletions
@@ -5,5 +5,5 @@ import com.sap.sailing.server.replication.impl.ReplicationFactoryImpl;
public interface ReplicationFactory {
static ReplicationFactory INSTANCE = new ReplicationFactoryImpl();
ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, int servletPort, int jmsPort);
ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int jmsPort);
}
@@ -1,11 +1,10 @@
package com.sap.sailing.server.replication;
import java.io.IOException;
import java.net.MalformedURLException;
import java.net.URL;
import java.net.UnknownHostException;
import javax.jms.JMSException;
import javax.jms.TopicSubscriber;
import com.rabbitmq.client.QueueingConsumer;
/**
* Identifies a master server instance from which a replica can obtain an initial load and continuous updates.
@@ -19,11 +18,16 @@ public interface ReplicationMasterDescriptor {
URL getInitialLoadURL() throws MalformedURLException;
TopicSubscriber getTopicSubscriber(String clientID) throws JMSException, UnknownHostException;
int getJMSPort();
int getServletPort();
String getHostname();
/**
* Creates a queue, declares the master's fanout exchange on the calling client and binds the queue to the exchange.
* Then, adds a consumer to the queue just created and starts consuming. The caller may keep calling
* {@link QueueingConsumer#nextDelivery()} on the consumer returned in order to obtain the next message.
*/
QueueingConsumer getConsumer() throws IOException;
}
@@ -3,8 +3,6 @@ package com.sap.sailing.server.replication;
import java.io.IOException;
import java.util.Map;
import javax.jms.JMSException;
import com.sap.sailing.server.RacingEventServiceOperation;
import com.sap.sailing.server.replication.impl.ReplicationServlet;
@@ -27,16 +25,15 @@ public interface ReplicationService {
* the JMS replication topic is created, then subscribing for the master's JMS replication topic and asking the servlet
* for the stream containing the initial load.
*/
void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException, JMSException;
void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException;
/**
* Registers a replica with this master instance. If the replication topic hasn't been created in the
* JMS message broker yet, it will be when this method returns. The <code>replica</code> will be considered
* in the result of {@link #getReplicaInfo()} when this call has succeeded.
* Registers a replica with this master instance. The <code>replica</code> will be considered in the result of
* {@link #getReplicaInfo()} when this call has succeeded.
*/
void registerReplica(ReplicaDescriptor replica) throws JMSException;
void registerReplica(ReplicaDescriptor replica);
void unregisterReplica(ReplicaDescriptor replica) throws JMSException;
void unregisterReplica(ReplicaDescriptor replica) throws IOException;
/**
* For a replica replicating off this master, provides statistics in the form of number of operations sent to that
@@ -1,7 +1,5 @@
package com.sap.sailing.server.replication.impl;
import java.io.File;
import java.io.FileNotFoundException;
import java.util.logging.Logger;
import org.osgi.framework.BundleActivator;
@@ -12,72 +10,25 @@ import com.sap.sailing.server.replication.ReplicationService;
public class Activator implements BundleActivator {
private static final Logger logger = Logger.getLogger(Activator.class.getName());
private static final String REPLICATION_PERSISTENCE_DIR_PROPERTY = "replication.persistenceDir";
private static final String BROKER_URL_PROPERTY = "replication.brokerURL";
private static final String REPLICATION_USE_JMX_PROPERTY = "replication.useJMX";
private MessageBrokerManager messageBrokerManager;
private static final String PROPERTY_NAME_EXCHANGE_NAME = "replication.exchangeName";
private ReplicationInstancesManager replicationInstancesManager;
private static BundleContext defaultContext;
public void start(BundleContext bundleContext) throws Exception {
defaultContext = bundleContext;
String replicationPersistenceDirectory = bundleContext.getProperty(REPLICATION_PERSISTENCE_DIR_PROPERTY);
if (replicationPersistenceDirectory == null) {
replicationPersistenceDirectory = System.getProperty("java.io.tmpdir");
String exchangeName = bundleContext.getProperty(PROPERTY_NAME_EXCHANGE_NAME);
if (exchangeName == null) {
exchangeName = "sapsailinganalytics";
}
String brokerURL = bundleContext.getProperty(BROKER_URL_PROPERTY);
if (brokerURL == null) {
brokerURL = "tcp://localhost:61616";
}
final File brokerPersistenceDir = new File(replicationPersistenceDirectory, "kahadb");
removeTemporaryTestBrokerPersistenceDirectory(brokerPersistenceDir);
MessageBrokerConfiguration brokerConfig = new MessageBrokerConfiguration("SailingServerReplicationBroker",
brokerURL, brokerPersistenceDir.getAbsolutePath());
messageBrokerManager = new MessageBrokerManager(brokerConfig);
String useJMX = bundleContext.getProperty(REPLICATION_USE_JMX_PROPERTY);
if (useJMX == null || useJMX.length() == 0) {
useJMX = "false";
}
messageBrokerManager.startMessageBroker(Boolean.valueOf(useJMX));
messageBrokerManager.createAndStartConnection();
replicationInstancesManager = new ReplicationInstancesManager();
ReplicationService serverReplicationMasterService = new ReplicationServiceImpl(replicationInstancesManager, messageBrokerManager);
ReplicationService serverReplicationMasterService = new ReplicationServiceImpl(exchangeName, replicationInstancesManager);
bundleContext.registerService(ReplicationService.class, serverReplicationMasterService, null);
logger.info("Registered replication service "+serverReplicationMasterService);
}
public static void removeTemporaryTestBrokerPersistenceDirectory(File brokerPersistenceDir) throws FileNotFoundException {
if (brokerPersistenceDir.exists() && brokerPersistenceDir.isDirectory()) {
logger.info("Deleting message broker persistence director "+brokerPersistenceDir);
deleteRecursive(brokerPersistenceDir);
}
File failoverStore = new File("activemq-data");
if (failoverStore.exists() && failoverStore.isDirectory()) {
logger.info("Deleting message broker failover store "+failoverStore);
deleteRecursive(failoverStore);
}
}
public static boolean deleteRecursive(File path) throws FileNotFoundException{
if (!path.exists()) throw new FileNotFoundException(path.getAbsolutePath());
boolean ret = true;
if (path.isDirectory()){
for (File f : path.listFiles()){
ret = ret && deleteRecursive(f);
}
}
return ret && path.delete();
}
public void stop(BundleContext bundleContext) throws Exception {
messageBrokerManager.closeSessions();
messageBrokerManager.closeConnections();
messageBrokerManager.stopMessageBroker();
}
public static BundleContext getDefaultContext() {
@@ -1,32 +0,0 @@
package com.sap.sailing.server.replication.impl;
/**
* A simple configuration for the message broker
*/
public class MessageBrokerConfiguration {
private final String brokerName;
private final String brokerUrl;
private final String dataStoreDirectory;
public MessageBrokerConfiguration(String brokerName, String brokerUrl, String dataStoreDirectory) {
super();
this.brokerName = brokerName;
this.brokerUrl = brokerUrl;
this.dataStoreDirectory = dataStoreDirectory;
}
public String getBrokerName() {
return brokerName;
}
public String getDataStoreDirectory() {
return dataStoreDirectory;
}
public String getBrokerUrl() {
return brokerUrl;
}
}
@@ -1,73 +0,0 @@
package com.sap.sailing.server.replication.impl;
import java.net.URI;
import javax.jms.Connection;
import javax.jms.JMSException;
import javax.jms.Session;
import org.apache.activemq.ActiveMQConnection;
import org.apache.activemq.ActiveMQConnectionFactory;
import org.apache.activemq.broker.BrokerService;
import org.apache.activemq.broker.TransportConnector;
public class MessageBrokerManager {
private final MessageBrokerConfiguration configuration;
private ActiveMQConnectionFactory connectionFactory;
private Connection connection;
private Session session;
private BrokerService broker;
public MessageBrokerManager(final MessageBrokerConfiguration configuration) {
this.configuration = configuration;
}
public void startMessageBroker(boolean useJmx) throws Exception {
broker = new BrokerService();
broker.setBrokerName(configuration.getBrokerName());
if (configuration.getDataStoreDirectory() != null) {
broker.setDataDirectory(configuration.getDataStoreDirectory());
}
broker.setUseJmx(useJmx);
TransportConnector transportConnector = new TransportConnector();
transportConnector.setUri(new URI(configuration.getBrokerUrl()));
broker.addConnector(transportConnector);
broker.start();
}
public void stopMessageBroker() throws Exception {
if (broker != null) {
broker.stop();
}
}
public void createAndStartConnection() throws JMSException {
connectionFactory = new ActiveMQConnectionFactory(ActiveMQConnection.DEFAULT_USER,
ActiveMQConnection.DEFAULT_PASSWORD, configuration.getBrokerUrl());
connection = connectionFactory.createConnection();
connection.start();
}
public void closeConnections() throws JMSException {
if (connection != null) {
connection.close();
}
}
public Session createSession(boolean transacted) throws JMSException {
session = connection.createSession(transacted, Session.AUTO_ACKNOWLEDGE);
return session;
}
public void closeSessions() throws JMSException {
if (session != null) {
session.close();
}
}
public Session getSession() {
return session;
}
}
@@ -4,10 +4,9 @@ import com.sap.sailing.server.replication.ReplicationFactory;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
public class ReplicationFactoryImpl implements ReplicationFactory {
@Override
public ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, int servletPort, int jmsPort) {
return new ReplicationMasterDescriptorImpl(hostname, servletPort, jmsPort);
public ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int messagingPort) {
return new ReplicationMasterDescriptorImpl(hostname, exchangeName, servletPort, messagingPort);
}
}
@@ -1,32 +1,27 @@
package com.sap.sailing.server.replication.impl;
import java.net.InetAddress;
import java.io.IOException;
import java.net.MalformedURLException;
import java.net.URL;
import java.net.UnknownHostException;
import javax.jms.Connection;
import javax.jms.JMSException;
import javax.jms.Session;
import javax.jms.Topic;
import javax.jms.TopicSubscriber;
import org.apache.activemq.ActiveMQConnection;
import org.apache.activemq.ActiveMQConnectionFactory;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.QueueingConsumer;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
import com.sap.sailing.server.replication.ReplicationService;
public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescriptor {
private static final String REPLICATION_SERVLET = "/replication/replication";
private final String hostname;
private final String exchangeName;
private final int servletPort;
private final int jmsPort;
public ReplicationMasterDescriptorImpl(String hostname, int servletPort, int jmsPort) {
public ReplicationMasterDescriptorImpl(String hostname, String exchangeName, int servletPort, int jmsPort) {
this.hostname = hostname;
this.servletPort = servletPort;
this.jmsPort = jmsPort;
this.exchangeName = exchangeName;
}
@Override
@@ -42,15 +37,17 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
}
@Override
public TopicSubscriber getTopicSubscriber(String clientID) throws JMSException, UnknownHostException {
ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(ActiveMQConnection.DEFAULT_USER,
ActiveMQConnection.DEFAULT_PASSWORD, "tcp://" + hostname + ":" + jmsPort);
connectionFactory.setClientID(clientID);
Connection connection = connectionFactory.createConnection();
connection.start();
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
Topic topic = session.createTopic(ReplicationService.SAILING_SERVER_REPLICATION_TOPIC);
return session.createDurableSubscriber(topic, InetAddress.getLocalHost().getHostAddress());
public QueueingConsumer getConsumer() throws IOException {
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost(getHostname());
Connection connection = connectionFactory.newConnection();
Channel channel = connection.createChannel();
channel.exchangeDeclare(exchangeName, "fanout");
QueueingConsumer consumer = new QueueingConsumer(channel);
String queueName = channel.queueDeclare().getQueue();
channel.queueBind(queueName, exchangeName, "");
channel.basicConsume(queueName, consumer);
return consumer;
}
@Override
@@ -11,16 +11,12 @@ import java.net.URLConnection;
import java.util.HashMap;
import java.util.Map;
import javax.jms.BytesMessage;
import javax.jms.DeliveryMode;
import javax.jms.JMSException;
import javax.jms.MessageProducer;
import javax.jms.Session;
import javax.jms.Topic;
import javax.jms.TopicSubscriber;
import org.osgi.util.tracker.ServiceTracker;
import com.rabbitmq.client.AMQP.Exchange;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.QueueingConsumer;
import com.sap.sailing.server.OperationExecutionListener;
import com.sap.sailing.server.RacingEventService;
import com.sap.sailing.server.RacingEventServiceOperation;
@@ -31,8 +27,8 @@ import com.sap.sailing.server.replication.ReplicationService;
/**
* Can observe a {@link RacingEventService} for the operations it performs that require replication. Only observes as
* long as there are replicas registered. If the last replica is de-registered, the service stops observing the
* {@link RacingEventService}. Operations received that require replication are broadcast to the
* {@link #getReplicationTopic() replication topic}.
* {@link RacingEventService}. Operations received that require replication are sent to the {@link Exchange} to which
* replica queues can bind. The exchange name is provided to this service during construction.
* <p>
*
* This service object {@link RacingEventService#addOperationExecutionListener(OperationExecutionListener) registers} as
@@ -46,12 +42,6 @@ import com.sap.sailing.server.replication.ReplicationService;
public class ReplicationServiceImpl implements ReplicationService, OperationExecutionListener, HasRacingEventService {
private final ReplicationInstancesManager replicationInstancesManager;
private final MessageBrokerManager messageBrokerManager;
private MessageProducer messageProducer;
private Topic replicationTopic;
private ServiceTracker<RacingEventService, RacingEventService> racingEventServiceTracker;
private final RacingEventService localService;
@@ -66,29 +56,43 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
*/
private final Map<ReplicationMasterDescriptor, String> replicaUUIDs;
public ReplicationServiceImpl(final ReplicationInstancesManager replicationInstancesManager,
final MessageBrokerManager messageBrokerManager) throws Exception {
private final Channel channel;
/**
* The name of the RabbitMQ exchange to which this replication service sends its replication operations in
* serialized form. Clients need to know this name to be able to bind their queues to the exchange.
*/
private final String exchangeName;
public ReplicationServiceImpl(String exchangeName, final ReplicationInstancesManager replicationInstancesManager) throws IOException {
this.replicationInstancesManager = replicationInstancesManager;
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>();
this.messageBrokerManager = messageBrokerManager;
racingEventServiceTracker = new ServiceTracker<RacingEventService, RacingEventService>(
Activator.getDefaultContext(), RacingEventService.class.getName(), null);
racingEventServiceTracker.open();
localService = null;
this.exchangeName = exchangeName;
channel = createChannel(exchangeName);
}
/**
* Like {@link #ReplicationServiceImpl(ReplicationInstancesManager, MessageBrokerManager)}, only that instead of using
* Like {@link #ReplicationServiceImpl(String, ReplicationInstancesManager)}, only that instead of using
* an OSGi service tracker to discover the {@link RacingEventService}, the service to replicate is "injected" here.
* @param exchangeName the name of the exchange to which replicas can bind
*/
public ReplicationServiceImpl(final ReplicationInstancesManager replicationInstancesManager,
final MessageBrokerManager messageBrokerManager, RacingEventService localService) {
public ReplicationServiceImpl(String exchangeName,
final ReplicationInstancesManager replicationInstancesManager, RacingEventService localService) throws IOException {
this.replicationInstancesManager = replicationInstancesManager;
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>();
this.messageBrokerManager = messageBrokerManager;
this.localService = localService;
this.exchangeName = exchangeName;
channel = createChannel(exchangeName);
}
private Channel createChannel(String exchangeName) throws IOException {
return new ConnectionFactory().newConnection().createChannel();
}
@Override
public RacingEventService getRacingEventService() {
RacingEventService result;
@@ -101,13 +105,9 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
}
@Override
public void registerReplica(ReplicaDescriptor replica) throws JMSException {
Topic topic = getReplicationTopic();
assert topic != null;
public void registerReplica(ReplicaDescriptor replica) {
if (!replicationInstancesManager.hasReplicas()) {
addAsListenerToRacingEventService();
messageBrokerManager.createAndStartConnection();
messageBrokerManager.createSession(/* transacted */ false);
}
replicationInstancesManager.registerReplica(replica);
}
@@ -117,53 +117,27 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
}
@Override
public void unregisterReplica(ReplicaDescriptor replica) throws JMSException {
public void unregisterReplica(ReplicaDescriptor replica) throws IOException {
replicationInstancesManager.unregisterReplica(replica);
if (!replicationInstancesManager.hasReplicas()) {
removeAsListenerFromRacingEventService();
messageBrokerManager.closeSessions();
messageBrokerManager.closeConnections();
}
}
private void removeAsListenerFromRacingEventService() {
getRacingEventService().removeOperationExecutionListener(this);
messageProducer = null;
}
private Topic getReplicationTopic() throws JMSException{
if (replicationTopic == null) {
Session session = messageBrokerManager.getSession();
if (session == null) {
session = messageBrokerManager.createSession(true);
}
replicationTopic = session.createTopic(SAILING_SERVER_REPLICATION_TOPIC);
}
return replicationTopic;
}
private void broadcastOperation(RacingEventServiceOperation<?> operation) throws Exception {
Topic topic = getReplicationTopic();
Session session = messageBrokerManager.getSession();
getMessageProducer(topic).setDeliveryMode(DeliveryMode.NON_PERSISTENT);
BytesMessage operationAsMessage = session.createBytesMessage();
// serialize operation into message
ByteArrayOutputStream bos = new ByteArrayOutputStream();
ObjectOutputStream oos = new ObjectOutputStream(bos);
oos.writeObject(operation);
oos.close();
operationAsMessage.writeBytes(bos.toByteArray());
messageProducer.send(operationAsMessage);
channel.basicPublish(exchangeName, /* routingKey */ "", /* properties */ null, bos.toByteArray());
replicationInstancesManager.log(operation);
}
private MessageProducer getMessageProducer(Topic topic) throws JMSException {
if (messageProducer == null) {
messageProducer = messageBrokerManager.getSession().createProducer(topic);
}
return messageProducer;
}
@Override
public Iterable<ReplicaDescriptor> getReplicaInfo() {
return replicationInstancesManager.getReplicaDescriptors();
@@ -175,13 +149,12 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
}
@Override
public void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException, JMSException {
public void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException {
replicatingFromMaster = master;
String uuid = registerReplicaWithMaster(master);
TopicSubscriber replicationSubscription = master.getTopicSubscriber(uuid);
registerReplicaWithMaster(master);
QueueingConsumer consumer = master.getConsumer();
URL initialLoadURL = master.getInitialLoadURL();
final Replicator replicator = new Replicator(master, this, /* startSuspended */ true);
replicationSubscription.setMessageListener(replicator);
final Replicator replicator = new Replicator(master, this, /* startSuspended */ true, consumer);
InputStream is = initialLoadURL.openStream();
ObjectInputStream ois = new ObjectInputStream(is) {
@Override
@@ -7,7 +7,6 @@ import java.net.UnknownHostException;
import java.util.Arrays;
import java.util.logging.Logger;
import javax.jms.JMSException;
import javax.servlet.ServletException;
import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse;
@@ -17,8 +16,8 @@ import org.osgi.util.tracker.ServiceTracker;
import com.sap.sailing.server.RacingEventService;
import com.sap.sailing.server.SailingServerHttpServlet;
import com.sap.sailing.server.replication.ReplicationService;
import com.sap.sailing.server.replication.ReplicaDescriptor;
import com.sap.sailing.server.replication.ReplicationService;
/**
* As the response to any type of <code>GET</code> request, sends a serialized copy of the {@link RacingEventService} to
@@ -60,11 +59,7 @@ public class ReplicationServlet extends SailingServerHttpServlet {
String action = req.getParameter(ACTION);
switch (Action.valueOf(action)) {
case REGISTER:
try {
registerClientWithReplicationService(req, resp);
} catch (JMSException e) {
resp.sendError(HttpServletResponse.SC_INTERNAL_SERVER_ERROR, e.getMessage());
}
registerClientWithReplicationService(req, resp);
break;
case INITIAL_LOAD:
ObjectOutputStream oos = new ObjectOutputStream(resp.getOutputStream());
@@ -84,7 +79,7 @@ public class ReplicationServlet extends SailingServerHttpServlet {
}
private void registerClientWithReplicationService(HttpServletRequest req, HttpServletResponse resp)
throws JMSException, IOException {
throws IOException {
ReplicaDescriptor replica = getReplicaDescriptor(req);
getReplicationService().registerReplica(replica);
resp.setContentType("text/plain");
@@ -6,12 +6,12 @@ import java.io.ObjectInputStream;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List;
import java.util.logging.Logger;
import javax.jms.BytesMessage;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageListener;
import com.rabbitmq.client.ConsumerCancelledException;
import com.rabbitmq.client.QueueingConsumer;
import com.rabbitmq.client.QueueingConsumer.Delivery;
import com.rabbitmq.client.ShutdownSignalException;
import com.sap.sailing.domain.base.DomainFactory;
import com.sap.sailing.server.RacingEventService;
import com.sap.sailing.server.RacingEventServiceOperation;
@@ -30,10 +30,13 @@ import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
* @author Axel Uhl (d043530)
*
*/
public class Replicator implements MessageListener {
public class Replicator implements Runnable {
private final static Logger logger = Logger.getLogger(Replicator.class.getName());
private final ReplicationMasterDescriptor master;
private final HasRacingEventService racingEventServiceTracker;
private final List<RacingEventServiceOperation<?>> queue;
private final QueueingConsumer consumer;
/**
* If the replicator is suspended, messages received are queued.
@@ -47,46 +50,52 @@ public class Replicator implements MessageListener {
* descriptor of the master server from which this replicator receives messages
* @param racingEventServiceTracker
* OSGi service tracker for the replica to which to apply the messages received
* @param consumer the RabbitMQ consumer from which to load messages
*/
public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker) {
this(master, racingEventServiceTracker, /* startSuspended */ false);
public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker, QueueingConsumer consumer) {
this(master, racingEventServiceTracker, /* startSuspended */ false, consumer);
}
public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker, boolean startSuspended) {
public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker, boolean startSuspended, QueueingConsumer consumer) {
this.queue = new ArrayList<RacingEventServiceOperation<?>>();
this.master = master;
this.racingEventServiceTracker = racingEventServiceTracker;
this.suspended = startSuspended;
this.consumer = consumer;
}
/**
* Starts fetching messages from the {@link #consumer}. After receiving a single message, assumes it's a serialized
* {@link RacingEventServiceOperation}, and applies it to the {@link RacingEventService} which is obtained from the
* service tracker passed to this replicator at construction time.
*/
@Override
public void run() {
ClassLoader oldClassLoader = Thread.currentThread().getContextClassLoader();
while (true) {
try {
Delivery delivery = consumer.nextDelivery();
byte[] bytesFromMessage = delivery.getBody();
// 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());
ObjectInputStream ois = DomainFactory.INSTANCE.createObjectInputStreamResolvingAgainstThisFactory(
new ByteArrayInputStream(bytesFromMessage));
RacingEventServiceOperation<?> operation = (RacingEventServiceOperation<?>) ois.readObject();
applyOrQueue(operation);
} catch (ShutdownSignalException | ConsumerCancelledException | InterruptedException | IOException | ClassNotFoundException e) {
logger.info("Exception while processing replica: "+e.getMessage());
logger.throwing(Replicator.class.getName(), "run", e);
} finally {
Thread.currentThread().setContextClassLoader(oldClassLoader);
}
}
}
public synchronized boolean isQueueEmpty() {
return queue.isEmpty();
}
/**
* Receives a single message, assuming it's a {@link RacingEventServiceOperation}, and applies it to the
* {@link RacingEventService} which is obtained from the service tracker passed to this replicator at construction
* time.
*/
@Override
public synchronized void onMessage(Message m) {
ClassLoader oldClassLoader = Thread.currentThread().getContextClassLoader();
try {
byte[] bytesFromMessage = getBytes((BytesMessage) m);
// 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());
ObjectInputStream ois = DomainFactory.INSTANCE.createObjectInputStreamResolvingAgainstThisFactory(
new ByteArrayInputStream(bytesFromMessage));
RacingEventServiceOperation<?> operation = (RacingEventServiceOperation<?>) ois.readObject();
applyOrQueue(operation);
} catch (IOException | ClassNotFoundException | JMSException e) {
throw new RuntimeException(e);
} finally {
Thread.currentThread().setContextClassLoader(oldClassLoader);
}
}
/**
* If the replicator is currently {@link #suspended}, the <code>operation</code> is queued, otherwise immediately applied to
* the receiving replica.
@@ -134,12 +143,6 @@ public class Replicator implements MessageListener {
return suspended;
}
private byte[] getBytes(BytesMessage m) throws JMSException {
byte[] buf = new byte[(int) m.getBodyLength()];
m.readBytes(buf);
return buf;
}
@Override
public String toString() {
return "Replicator for master "+master+", queue size: "+queue.size();
@@ -1,27 +0,0 @@
package com.sap.sailing.server.replication.impl;
import javax.jms.ExceptionListener;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageListener;
import javax.jms.TextMessage;
public class SampleMessageConsumer implements MessageListener, ExceptionListener {
synchronized public void onException(JMSException ex) {
System.out.println("JMS Exception occured: " + ex.getMessage());
}
public void onMessage(Message message) {
if (message instanceof TextMessage) {
TextMessage textMessage = (TextMessage) message;
try {
System.out.println("Received message: " + textMessage.getText());
} catch (JMSException ex) {
System.out.println("Error reading message: " + ex);
}
} else {
System.out.println("Received: " + message);
}
}
}