mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-09-30 09:26:44 +00:00
replaced jms queue with topic for the operation replication
This commit is contained in:
+8
-8
@@ -5,9 +5,9 @@ import java.io.Serializable;
|
||||
import javax.jms.DeliveryMode;
|
||||
import javax.jms.JMSException;
|
||||
import javax.jms.MessageProducer;
|
||||
import javax.jms.Queue;
|
||||
import javax.jms.Session;
|
||||
import javax.jms.TextMessage;
|
||||
import javax.jms.Topic;
|
||||
|
||||
import com.sap.sailing.server.operationaltransformation.RacingEventServiceOperation;
|
||||
import com.sap.sailing.server.replication.ServerReplicationMasterService;
|
||||
@@ -17,29 +17,29 @@ public class ServerReplicationMasterServiceImpl implements ServerReplicationMast
|
||||
|
||||
private final MessageBrokerManager messageBrokerManager;
|
||||
|
||||
private Queue replicationQueue;
|
||||
private Topic replicationTopic;
|
||||
|
||||
public ServerReplicationMasterServiceImpl(final ReplicationInstancesManager replicationInstancesManager, final MessageBrokerManager messageBrokerManager) {
|
||||
this.replicationInstancesManager = replicationInstancesManager;
|
||||
this.messageBrokerManager = messageBrokerManager;
|
||||
}
|
||||
|
||||
private Queue getReplicationQueue() throws JMSException{
|
||||
if(replicationQueue == null) {
|
||||
private Topic getReplicationTopic() throws JMSException{
|
||||
if(replicationTopic== null) {
|
||||
Session session = messageBrokerManager.getSession();
|
||||
if(session == null) {
|
||||
session = messageBrokerManager.createSession(true);
|
||||
}
|
||||
replicationQueue = session.createQueue("SailingServerManyPlayersQueue");
|
||||
replicationTopic = session.createTopic("SailingServerReplicationTopic");
|
||||
}
|
||||
return replicationQueue;
|
||||
return replicationTopic;
|
||||
}
|
||||
|
||||
public void broadcastOperation(RacingEventServiceOperation operation) throws Exception {
|
||||
Queue queue = getReplicationQueue();
|
||||
Topic topic = getReplicationTopic();
|
||||
Session session = messageBrokerManager.getSession();
|
||||
|
||||
MessageProducer producer = messageBrokerManager.getSession().createProducer(queue);
|
||||
MessageProducer producer = messageBrokerManager.getSession().createProducer(topic);
|
||||
producer.setDeliveryMode(DeliveryMode.PERSISTENT);
|
||||
TextMessage message = session.createTextMessage("Hello World!");
|
||||
System.out.println("Sending message: " + message.getText());
|
||||
|
||||
Reference in New Issue
Block a user