mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-10-11 06:40:53 +00:00
Now also the queue will survive a connection drop and will also handle
client disconnects gracefully and free contents after some time.
This commit is contained in:
1 parent
bf35cbf077
commit
464d7d2051
6 files changed
+50
-8
No files matched your search
+1
-1
@@ -5,5 +5,5 @@ import com.sap.sailing.server.replication.impl.ReplicationFactoryImpl;
|
||||
public interface ReplicationFactory {
|
||||
static ReplicationFactory INSTANCE = new ReplicationFactoryImpl();
|
||||
|
||||
ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int jmsPort);
|
||||
ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int jmsPort, String jmsQueueName);
|
||||
}
|
||||
+2
-2
@@ -5,8 +5,8 @@ import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
|
||||
|
||||
public class ReplicationFactoryImpl implements ReplicationFactory {
|
||||
@Override
|
||||
public ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int messagingPort) {
|
||||
return new ReplicationMasterDescriptorImpl(hostname, exchangeName, servletPort, messagingPort);
|
||||
public ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int messagingPort, String clientUUID) {
|
||||
return new ReplicationMasterDescriptorImpl(hostname, exchangeName, servletPort, messagingPort, clientUUID);
|
||||
}
|
||||
|
||||
}
|
||||
+44
-2
@@ -4,6 +4,8 @@ import java.io.IOException;
|
||||
import java.io.UnsupportedEncodingException;
|
||||
import java.net.MalformedURLException;
|
||||
import java.net.URL;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.UUID;
|
||||
|
||||
import com.rabbitmq.client.Channel;
|
||||
@@ -19,15 +21,17 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
|
||||
private final String exchangeName;
|
||||
private final int servletPort;
|
||||
private final int messagingPort;
|
||||
private final String queueName;
|
||||
|
||||
/**
|
||||
* @param messagingPort 0 means use default port
|
||||
*/
|
||||
public ReplicationMasterDescriptorImpl(String hostname, String exchangeName, int servletPort, int messagingPort) {
|
||||
public ReplicationMasterDescriptorImpl(String hostname, String exchangeName, int servletPort, int messagingPort, String queueName) {
|
||||
this.hostname = hostname;
|
||||
this.servletPort = servletPort;
|
||||
this.messagingPort = messagingPort;
|
||||
this.exchangeName = exchangeName;
|
||||
this.queueName = queueName;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -61,9 +65,47 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
|
||||
}
|
||||
Connection connection = connectionFactory.newConnection();
|
||||
Channel channel = connection.createChannel();
|
||||
|
||||
/*
|
||||
* Connect a queue to the given exchange that has already
|
||||
* been created by the master server.
|
||||
*/
|
||||
channel.exchangeDeclare(exchangeName, "fanout");
|
||||
QueueingConsumer consumer = new QueueingConsumer(channel);
|
||||
String queueName = channel.queueDeclare().getQueue();
|
||||
|
||||
/*
|
||||
* The x-message-ttl argument to queue.declare controls for how long a message published to a queue can live before
|
||||
* it is discarded. A message that has been in the queue for longer than the configured TTL is said to be dead.
|
||||
* Note that a message routed to multiple queues can die at different times, or not at all,
|
||||
* in each queue in which it resides. The death of a message in one queue has no impact on the life of the
|
||||
* same message in other queues.
|
||||
*/
|
||||
final Map<String, Object> args = new HashMap<String, Object>();
|
||||
args.put("x-message-ttl", (60*30)*1000); // messages will live half an hour in queue before being deleted
|
||||
|
||||
/*
|
||||
* The x-expires argument to queue.declare controls for how long a queue can be unused before it is automatically
|
||||
* deleted. Unused means the queue has no consumers, the queue has not been redeclared, and basic.get has not
|
||||
* been invoked for a duration of at least the expiration period.
|
||||
*/
|
||||
args.put("x-expires", (60*60)*1000); // queue will live one hour before being deleted
|
||||
|
||||
/*
|
||||
* The maximum length of a queue can be limited to a set number of messages by supplying the x-max-length queue
|
||||
* declaration argument with a non-negative integer value. Queue length is a measure that takes into account
|
||||
* ready messages, ignoring unacknowledged messages and message size. Messages will be dropped or dead-lettered
|
||||
* from the front of the queue to make room for new messages once the limit is reached.
|
||||
*/
|
||||
args.put("x-max-length", 3000000);
|
||||
|
||||
// a server-named non-exclusive, non-durable queue
|
||||
// this queue will survive a connection drop (autodelete=false) and
|
||||
// will also support being reconnected (exclusive=false). it will
|
||||
// not survive a rabbitmq server restart (durable=false).
|
||||
String queueName = channel.queueDeclare(this.queueName,
|
||||
/*durable*/ false, /*exclusive*/ false, /*auto-delete*/ false, args).getQueue();
|
||||
|
||||
// from now on we get all new messages that the exchange is getting from producer
|
||||
channel.queueBind(queueName, exchangeName, "");
|
||||
channel.basicConsume(queueName, /* auto-ack */ true, consumer);
|
||||
return consumer;
|
||||
|
||||
Reference in new issue
Block a user