mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-09-30 09:26:44 +00:00
support course change replication
This commit is contained in:
+10
-2
@@ -1,8 +1,10 @@
|
||||
package com.sap.sailing.server.replication.impl;
|
||||
|
||||
import java.io.ByteArrayOutputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.io.ObjectInputStream;
|
||||
import java.io.ObjectOutputStream;
|
||||
import java.net.URL;
|
||||
import java.net.URLConnection;
|
||||
import java.util.ArrayList;
|
||||
@@ -11,10 +13,10 @@ import java.util.Iterator;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import javax.jms.BytesMessage;
|
||||
import javax.jms.DeliveryMode;
|
||||
import javax.jms.JMSException;
|
||||
import javax.jms.MessageProducer;
|
||||
import javax.jms.ObjectMessage;
|
||||
import javax.jms.Session;
|
||||
import javax.jms.Topic;
|
||||
import javax.jms.TopicSubscriber;
|
||||
@@ -141,7 +143,13 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
|
||||
Topic topic = getReplicationTopic();
|
||||
Session session = messageBrokerManager.getSession();
|
||||
getMessageProducer(topic).setDeliveryMode(DeliveryMode.PERSISTENT);
|
||||
ObjectMessage operationAsMessage = session.createObjectMessage(operation);
|
||||
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);
|
||||
}
|
||||
|
||||
|
||||
+17
-3
@@ -1,10 +1,15 @@
|
||||
package com.sap.sailing.server.replication.impl;
|
||||
|
||||
import java.io.ByteArrayInputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.ObjectInputStream;
|
||||
|
||||
import javax.jms.BytesMessage;
|
||||
import javax.jms.JMSException;
|
||||
import javax.jms.Message;
|
||||
import javax.jms.MessageListener;
|
||||
import javax.jms.ObjectMessage;
|
||||
|
||||
import com.sap.sailing.domain.base.DomainFactory;
|
||||
import com.sap.sailing.server.RacingEventService;
|
||||
import com.sap.sailing.server.RacingEventServiceOperation;
|
||||
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
|
||||
@@ -34,13 +39,22 @@ public class Replicator implements MessageListener {
|
||||
@Override
|
||||
public void onMessage(Message m) {
|
||||
try {
|
||||
RacingEventServiceOperation<?> operation = (RacingEventServiceOperation<?>) ((ObjectMessage) m).getObject();
|
||||
byte[] bytesFromMessage = getBytes((BytesMessage) m);
|
||||
ObjectInputStream ois = DomainFactory.INSTANCE.createObjectInputStreamResolvingAgainstThisFactory(
|
||||
new ByteArrayInputStream(bytesFromMessage));
|
||||
RacingEventServiceOperation<?> operation = (RacingEventServiceOperation<?>) ois.readObject();
|
||||
racingEventServiceTracker.getRacingEventService().apply(operation);
|
||||
} catch (JMSException e) {
|
||||
} catch (IOException | ClassNotFoundException | JMSException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
|
||||
Reference in New Issue
Block a user