mirror of
https://github.com/ebean-orm/ebean.git
synced 2024-04-21 10:51:47 +00:00
Refactor properties bootup, remove GlobalProperties and allow ServerConfig.loadFromProperties()
This commit is contained in:
@@ -17,9 +17,6 @@ import java.io.DataOutputStream;
|
||||
* be common for many Ack, Resend and Control messages to all be contained in a
|
||||
* single packet.
|
||||
* </p>
|
||||
*
|
||||
* @author rbygrave
|
||||
*
|
||||
*/
|
||||
public class BinaryMessage {
|
||||
|
||||
@@ -28,8 +25,7 @@ public class BinaryMessage {
|
||||
public static final int TYPE_TABLEIUD = 2;
|
||||
public static final int TYPE_BEANDELTA = 3;
|
||||
public static final int TYPE_BEANPATHUPDATE = 4;
|
||||
public static final int TYPE_INDEX_INVALIDATE = 6;
|
||||
public static final int TYPE_INDEX = 7;
|
||||
|
||||
public static final int TYPE_MSGACK = 8;
|
||||
public static final int TYPE_MSGRESEND = 9;
|
||||
|
||||
|
||||
@@ -3,8 +3,7 @@ package com.avaje.ebeaninternal.server.cluster;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
import com.avaje.ebean.EbeanServer;
|
||||
import com.avaje.ebean.config.GlobalProperties;
|
||||
import com.avaje.ebeaninternal.api.ClassUtil;
|
||||
import com.avaje.ebean.config.ContainerConfig;
|
||||
import com.avaje.ebeaninternal.server.cluster.mcast.McastClusterManager;
|
||||
import com.avaje.ebeaninternal.server.cluster.socket.SocketClusterBroadcast;
|
||||
import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent;
|
||||
@@ -26,41 +25,36 @@ public class ClusterManager {
|
||||
|
||||
private boolean started;
|
||||
|
||||
public ClusterManager() {
|
||||
public ClusterManager(ContainerConfig containerConfig) {
|
||||
|
||||
String clusterType = GlobalProperties.get("ebean.cluster.type", null);
|
||||
if (clusterType == null || clusterType.trim().length() == 0) {
|
||||
// not clustering this instance
|
||||
this.broadcast = null;
|
||||
|
||||
} else {
|
||||
|
||||
try {
|
||||
if ("mcast".equalsIgnoreCase(clusterType)) {
|
||||
this.broadcast = new McastClusterManager();
|
||||
|
||||
} else if ("socket".equalsIgnoreCase(clusterType)) {
|
||||
this.broadcast = new SocketClusterBroadcast();
|
||||
|
||||
} else {
|
||||
logger.info("Clustering using [" + clusterType + "]");
|
||||
this.broadcast = (ClusterBroadcast) ClassUtil.newInstance(clusterType);
|
||||
ContainerConfig.ClusterMode mode = containerConfig.getMode();
|
||||
try {
|
||||
switch (mode) {
|
||||
case SOCKET: {
|
||||
this.broadcast = new SocketClusterBroadcast(containerConfig);
|
||||
break;
|
||||
}
|
||||
case MULTICAST: {
|
||||
this.broadcast = new McastClusterManager(containerConfig);
|
||||
break;
|
||||
}
|
||||
default: {
|
||||
this.broadcast = null;
|
||||
}
|
||||
|
||||
} catch (Exception e) {
|
||||
String msg = "Error initialising ClusterManager type [" + clusterType + "]";
|
||||
logger.error(msg, e);
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
|
||||
} catch (Exception e) {
|
||||
logger.error("Error initialising ClusterManager type [" + mode + "]", e);
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
|
||||
public void registerServer(EbeanServer server) {
|
||||
synchronized (monitor) {
|
||||
serverMap.put(server.getName(), server);
|
||||
if (!started) {
|
||||
startup();
|
||||
}
|
||||
serverMap.put(server.getName(), server);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -7,7 +7,6 @@ import com.avaje.ebeaninternal.api.SpiEbeanServer;
|
||||
import com.avaje.ebeaninternal.api.TransactionEventTable.TableIUD;
|
||||
import com.avaje.ebeaninternal.server.transaction.BeanDelta;
|
||||
import com.avaje.ebeaninternal.server.transaction.BeanPersistIds;
|
||||
import com.avaje.ebeaninternal.server.transaction.IndexEvent;
|
||||
import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent;
|
||||
|
||||
/**
|
||||
@@ -16,7 +15,6 @@ import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent;
|
||||
* Due to the hard limit for UDP packet sizes a RemoteTransactionEvent
|
||||
* is actually broken up into smaller messages.
|
||||
* </p>
|
||||
* @author rbygrave
|
||||
*/
|
||||
public class PacketTransactionEvent extends Packet {
|
||||
|
||||
@@ -63,10 +61,6 @@ public class PacketTransactionEvent extends Packet {
|
||||
event.addBeanDelta(BeanDelta.readBinaryMessage(server, dataInput));
|
||||
break;
|
||||
|
||||
case BinaryMessage.TYPE_INDEX:
|
||||
event.addIndexEvent(IndexEvent.readBinaryMessage(dataInput));
|
||||
break;
|
||||
|
||||
default:
|
||||
throw new RuntimeException("Invalid Transaction msgType "+msgType);
|
||||
}
|
||||
|
||||
+41
-36
@@ -7,49 +7,54 @@ import java.util.List;
|
||||
|
||||
import com.avaje.ebeaninternal.api.SpiEbeanServer;
|
||||
import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* Mechanism to convert RemoteTransactionEvent to/from byte[] content.
|
||||
*/
|
||||
public abstract class SerialiseTransactionHelper {
|
||||
|
||||
private final PacketWriter packetWriter;
|
||||
private static final Logger logger = LoggerFactory.getLogger(SerialiseTransactionHelper.class);
|
||||
|
||||
public SerialiseTransactionHelper() {
|
||||
packetWriter = new PacketWriter(Integer.MAX_VALUE);
|
||||
private final PacketWriter packetWriter;
|
||||
|
||||
public SerialiseTransactionHelper() {
|
||||
packetWriter = new PacketWriter(Integer.MAX_VALUE);
|
||||
}
|
||||
|
||||
public abstract SpiEbeanServer getEbeanServer(String serverName);
|
||||
|
||||
/**
|
||||
* Convert the RemoteTransactionEvent to byte[] content.
|
||||
*/
|
||||
public DataHolder createDataHolder(RemoteTransactionEvent transEvent) throws IOException {
|
||||
|
||||
List<Packet> packetList = packetWriter.write(transEvent);
|
||||
if (packetList.size() != 1) {
|
||||
throw new RuntimeException("Always expecting 1 Packet but got " + packetList.size());
|
||||
}
|
||||
byte[] data = packetList.get(0).getBytes();
|
||||
return new DataHolder(data);
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert the byte[] content to RemoteTransactionEvent.
|
||||
*/
|
||||
public RemoteTransactionEvent read(DataHolder dataHolder) throws IOException {
|
||||
|
||||
ByteArrayInputStream bi = new ByteArrayInputStream(dataHolder.getData());
|
||||
DataInputStream dataInput = new DataInputStream(bi);
|
||||
|
||||
Packet header = Packet.readHeader(dataInput);
|
||||
|
||||
SpiEbeanServer server = getEbeanServer(header.getServerName());
|
||||
if (server == null) {
|
||||
logger.error("server [{}] not found/registered?", header.getServerName());
|
||||
}
|
||||
|
||||
public abstract SpiEbeanServer getEbeanServer(String serverName);
|
||||
|
||||
/**
|
||||
* Convert the RemoteTransactionEvent to byte[] content.
|
||||
*/
|
||||
public DataHolder createDataHolder(RemoteTransactionEvent transEvent) throws IOException {
|
||||
|
||||
List<Packet> packetList = packetWriter.write(transEvent);
|
||||
if (packetList.size() != 1) {
|
||||
throw new RuntimeException("Always expecting 1 Packet but got " + packetList.size());
|
||||
}
|
||||
byte[] data = packetList.get(0).getBytes();
|
||||
return new DataHolder(data);
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert the byte[] content to RemoteTransactionEvent.
|
||||
*/
|
||||
public RemoteTransactionEvent read(DataHolder dataHolder) throws IOException {
|
||||
|
||||
ByteArrayInputStream bi = new ByteArrayInputStream(dataHolder.getData());
|
||||
DataInputStream dataInput = new DataInputStream(bi);
|
||||
|
||||
Packet header = Packet.readHeader(dataInput);
|
||||
|
||||
SpiEbeanServer server = getEbeanServer(header.getServerName());
|
||||
|
||||
PacketTransactionEvent tranEventPacket = PacketTransactionEvent.forRead(header, server);
|
||||
tranEventPacket.read(dataInput);
|
||||
|
||||
return tranEventPacket.getEvent();
|
||||
|
||||
}
|
||||
PacketTransactionEvent tranEventPacket = PacketTransactionEvent.forRead(header, server);
|
||||
tranEventPacket.read(dataInput);
|
||||
return tranEventPacket.getEvent();
|
||||
}
|
||||
}
|
||||
|
||||
+494
-515
File diff suppressed because it is too large
Load Diff
@@ -19,6 +19,12 @@
|
||||
*/
|
||||
package com.avaje.ebeaninternal.server.cluster.mcast;
|
||||
|
||||
import com.avaje.ebeaninternal.api.SpiEbeanServer;
|
||||
import com.avaje.ebeaninternal.server.cluster.Packet;
|
||||
import com.avaje.ebeaninternal.server.cluster.PacketTransactionEvent;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.ByteArrayInputStream;
|
||||
import java.io.DataInput;
|
||||
import java.io.DataInputStream;
|
||||
@@ -28,13 +34,6 @@ import java.net.InetAddress;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.net.MulticastSocket;
|
||||
|
||||
import com.avaje.ebean.config.GlobalProperties;
|
||||
import com.avaje.ebeaninternal.api.SpiEbeanServer;
|
||||
import com.avaje.ebeaninternal.server.cluster.Packet;
|
||||
import com.avaje.ebeaninternal.server.cluster.PacketTransactionEvent;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* Listens for Incoming packets.
|
||||
*
|
||||
@@ -56,8 +55,6 @@ public class McastListener implements Runnable {
|
||||
|
||||
private final InetAddress group;
|
||||
|
||||
private final boolean debugIgnore;
|
||||
|
||||
private DatagramPacket pack;
|
||||
|
||||
private byte[] receiveBuffer;
|
||||
@@ -73,8 +70,6 @@ public class McastListener implements Runnable {
|
||||
int bufferSize, int timeout, String localSenderHostPort,
|
||||
boolean disableLoopback, int ttl, InetAddress mcastBindAddress) {
|
||||
|
||||
this.debugIgnore = GlobalProperties.getBoolean("ebean.debug.mcast.ignore", false);
|
||||
|
||||
this.owner = owner;
|
||||
this.packetControl = packetControl;
|
||||
this.localSenderHostPort = localSenderHostPort;
|
||||
@@ -96,7 +91,7 @@ public class McastListener implements Runnable {
|
||||
this.sock.setSoTimeout(timeout);
|
||||
|
||||
if (disableLoopback){
|
||||
sock.setLoopbackMode(disableLoopback);
|
||||
sock.setLoopbackMode(true);
|
||||
}
|
||||
|
||||
if (mcastBindAddress != null) {
|
||||
@@ -173,7 +168,7 @@ public class McastListener implements Runnable {
|
||||
String senderHostPort = senderAddr.getAddress().getHostAddress()+":"+senderAddr.getPort();
|
||||
|
||||
if (senderHostPort.equals(localSenderHostPort)){
|
||||
if (debugIgnore || logger.isDebugEnabled()){
|
||||
if (logger.isTraceEnabled()){
|
||||
logger.info("Ignoring message as sent by localSender: "+localSenderHostPort);
|
||||
}
|
||||
} else {
|
||||
@@ -195,7 +190,7 @@ public class McastListener implements Runnable {
|
||||
boolean processThisPacket = ackMsg || packetControl.isProcessPacket(senderHostPort, header.getPacketId());
|
||||
|
||||
if (!processThisPacket){
|
||||
if (debugIgnore || logger.isDebugEnabled()){
|
||||
if (logger.isTraceEnabled()){
|
||||
logger.info("Already processed packet: "+header.getPacketId()+" type:"+header.getPacketType()+" len:"+data.length);
|
||||
}
|
||||
} else {
|
||||
|
||||
@@ -21,15 +21,18 @@ class RequestProcessor implements Runnable {
|
||||
private final Socket clientSocket;
|
||||
|
||||
private final SocketClusterBroadcast owner;
|
||||
|
||||
/**
|
||||
|
||||
private final String hostPort;
|
||||
|
||||
/**
|
||||
* Create including the Listener (used to lookup the Request Handler) and
|
||||
* the socket itself.
|
||||
*/
|
||||
public RequestProcessor(SocketClusterBroadcast owner, Socket clientSocket) {
|
||||
this.clientSocket = clientSocket;
|
||||
this.owner = owner;
|
||||
}
|
||||
this.hostPort = owner.getHostPort();
|
||||
}
|
||||
|
||||
/**
|
||||
* This will parse out the command. Lookup the appropriate Handler and
|
||||
@@ -39,22 +42,20 @@ class RequestProcessor implements Runnable {
|
||||
*/
|
||||
public void run() {
|
||||
try {
|
||||
SocketConnection sc = new SocketConnection(clientSocket);
|
||||
|
||||
while(true){
|
||||
if (owner.process(sc)) {
|
||||
// got the offline message or timeout
|
||||
break;
|
||||
}
|
||||
}
|
||||
sc.disconnect();
|
||||
|
||||
} catch (IOException e) {
|
||||
logger.error(null, e);
|
||||
} catch (ClassNotFoundException e) {
|
||||
logger.error(null, e);
|
||||
logger.trace("start listening for cluster messages");
|
||||
SocketConnection sc = new SocketConnection(clientSocket);
|
||||
while (true) {
|
||||
if (owner.process(sc)) {
|
||||
// got the offline message or timeout
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
logger.trace("disconnecting: {}", hostPort);
|
||||
sc.disconnect();
|
||||
|
||||
} catch (Exception e) {
|
||||
logger.error("Error listening for messages - "+owner.getHostPort(), e);
|
||||
}
|
||||
}
|
||||
|
||||
};
|
||||
}
|
||||
@@ -35,6 +35,10 @@ class SocketClient {
|
||||
this.hostPort = address.getHostName()+":"+address.getPort();
|
||||
}
|
||||
|
||||
public String toString() {
|
||||
return address.toString();
|
||||
}
|
||||
|
||||
public String getHostPort() {
|
||||
return hostPort;
|
||||
}
|
||||
|
||||
+217
-209
@@ -1,248 +1,256 @@
|
||||
package com.avaje.ebeaninternal.server.cluster.socket;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.InterruptedIOException;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.util.HashMap;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import javax.persistence.PersistenceException;
|
||||
|
||||
import com.avaje.ebean.config.GlobalProperties;
|
||||
import com.avaje.ebean.config.ContainerConfig;
|
||||
import com.avaje.ebeaninternal.api.SpiEbeanServer;
|
||||
import com.avaje.ebeaninternal.server.cluster.ClusterBroadcast;
|
||||
import com.avaje.ebeaninternal.server.cluster.ClusterManager;
|
||||
import com.avaje.ebeaninternal.server.cluster.DataHolder;
|
||||
import com.avaje.ebeaninternal.server.cluster.SerialiseTransactionHelper;
|
||||
import com.avaje.ebeaninternal.server.lib.util.StringHelper;
|
||||
import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import javax.persistence.PersistenceException;
|
||||
import java.io.EOFException;
|
||||
import java.io.IOException;
|
||||
import java.io.InterruptedIOException;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
/**
|
||||
* Broadcast messages across the cluster using sockets.
|
||||
* Broadcast messages across the cluster using sockets.
|
||||
*/
|
||||
public class SocketClusterBroadcast implements ClusterBroadcast {
|
||||
|
||||
private static final Logger logger = LoggerFactory.getLogger(SocketClusterBroadcast.class);
|
||||
|
||||
private final SocketClient local;
|
||||
|
||||
private final HashMap<String,SocketClient> clientMap;
|
||||
|
||||
private final SocketClusterListener listener;
|
||||
|
||||
private SocketClient[] members;
|
||||
private static final Logger logger = LoggerFactory.getLogger(SocketClusterBroadcast.class);
|
||||
|
||||
private ClusterManager clusterManager;
|
||||
private final SocketClient local;
|
||||
|
||||
private final TxnSerialiseHelper txnSerialiseHelper = new TxnSerialiseHelper();
|
||||
|
||||
private final AtomicInteger txnOutgoing = new AtomicInteger();
|
||||
private final AtomicInteger txnIncoming = new AtomicInteger();
|
||||
|
||||
|
||||
public SocketClusterBroadcast( ){
|
||||
|
||||
String localHostPort = GlobalProperties.get("ebean.cluster.local", null);
|
||||
String members = GlobalProperties.get("ebean.cluster.members", null);
|
||||
private final HashMap<String, SocketClient> clientMap;
|
||||
|
||||
logger.info("Clustering using Sockets local["+localHostPort+"] members["+members+"]");
|
||||
|
||||
this.local = new SocketClient(parseFullName(localHostPort));
|
||||
this.clientMap = new HashMap<String, SocketClient>();
|
||||
|
||||
String[] memArray = StringHelper.delimitedToArray(members, ",", false);
|
||||
for (int i = 0; i < memArray.length; i++) {
|
||||
InetSocketAddress member = parseFullName(memArray[i]);
|
||||
SocketClient client = new SocketClient(member);
|
||||
if (!local.getHostPort().equalsIgnoreCase(client.getHostPort())) {
|
||||
// don't add the local one ...
|
||||
clientMap.put(client.getHostPort(), client);
|
||||
}
|
||||
}
|
||||
|
||||
this.members = clientMap.values().toArray(new SocketClient[clientMap.size()]);
|
||||
this.listener = new SocketClusterListener(this, local.getPort());
|
||||
}
|
||||
|
||||
/**
|
||||
* Return the current status of this instance.
|
||||
*/
|
||||
public SocketClusterStatus getStatus() {
|
||||
|
||||
// count of online members
|
||||
int currentGroupSize = 0;
|
||||
for (int i = 0; i < members.length; i++) {
|
||||
if (members[i].isOnline()) {
|
||||
++currentGroupSize;
|
||||
}
|
||||
}
|
||||
int txnIn = txnIncoming.get();
|
||||
int txnOut = txnOutgoing.get();
|
||||
|
||||
return new SocketClusterStatus(currentGroupSize, txnIn, txnOut);
|
||||
}
|
||||
|
||||
public void startup(ClusterManager clusterManager) {
|
||||
|
||||
this.clusterManager = clusterManager;
|
||||
try {
|
||||
listener.startListening();
|
||||
register();
|
||||
private final SocketClusterListener listener;
|
||||
|
||||
} catch (IOException e) {
|
||||
throw new PersistenceException(e);
|
||||
}
|
||||
private SocketClient[] members;
|
||||
|
||||
private ClusterManager clusterManager;
|
||||
|
||||
private final TxnSerialiseHelper txnSerialiseHelper = new TxnSerialiseHelper();
|
||||
|
||||
private final AtomicInteger txnOutgoing = new AtomicInteger();
|
||||
private final AtomicInteger txnIncoming = new AtomicInteger();
|
||||
|
||||
public SocketClusterBroadcast(ContainerConfig containerConfig) {
|
||||
|
||||
ContainerConfig.SocketConfig socketConfig = containerConfig.getSocketConfig();
|
||||
|
||||
String localHostPort = socketConfig.getLocalHostPort();
|
||||
List<String> members = socketConfig.getMembers();
|
||||
|
||||
logger.info("Clustering using Sockets local[" + localHostPort + "] members[" + members + "]");
|
||||
|
||||
this.local = new SocketClient(parseFullName(localHostPort));
|
||||
this.clientMap = new HashMap<String, SocketClient>();
|
||||
|
||||
for (String memberHostPort : members) {
|
||||
InetSocketAddress member = parseFullName(memberHostPort);
|
||||
SocketClient client = new SocketClient(member);
|
||||
if (!local.getHostPort().equalsIgnoreCase(client.getHostPort())) {
|
||||
// don't add the local one ...
|
||||
clientMap.put(client.getHostPort(), client);
|
||||
}
|
||||
}
|
||||
|
||||
public void shutdown() {
|
||||
deregister();
|
||||
listener.shutdown();
|
||||
}
|
||||
|
||||
/**
|
||||
* Register with all the other members of the Cluster.
|
||||
*/
|
||||
private void register() {
|
||||
this.members = clientMap.values().toArray(new SocketClient[clientMap.size()]);
|
||||
this.listener = new SocketClusterListener(this, local.getPort(), socketConfig.getCoreThreads(), socketConfig.getMaxThreads(), socketConfig.getThreadPoolName());
|
||||
}
|
||||
|
||||
SocketClusterMessage h = SocketClusterMessage.register(local.getHostPort(), true);
|
||||
|
||||
for (int i = 0; i < members.length; i++) {
|
||||
boolean online = members[i].register(h);
|
||||
|
||||
String msg = "Cluster Member ["+members[i].getHostPort()+"] online["+online+"]";
|
||||
logger.info(msg);
|
||||
}
|
||||
public String getHostPort() {
|
||||
return local.getHostPort();
|
||||
}
|
||||
|
||||
/**
|
||||
* Return the current status of this instance.
|
||||
*/
|
||||
public SocketClusterStatus getStatus() {
|
||||
|
||||
// count of online members
|
||||
int currentGroupSize = 0;
|
||||
for (int i = 0; i < members.length; i++) {
|
||||
if (members[i].isOnline()) {
|
||||
++currentGroupSize;
|
||||
}
|
||||
}
|
||||
int txnIn = txnIncoming.get();
|
||||
int txnOut = txnOutgoing.get();
|
||||
|
||||
protected void setMemberOnline(String fullName, boolean online) throws IOException {
|
||||
synchronized (clientMap) {
|
||||
String msg = "Cluster Member ["+fullName+"] online["+online+"]";
|
||||
logger.info(msg);
|
||||
SocketClient member = clientMap.get(fullName);
|
||||
member.setOnline(online);
|
||||
}
|
||||
return new SocketClusterStatus(currentGroupSize, txnIn, txnOut);
|
||||
}
|
||||
|
||||
public void startup(ClusterManager clusterManager) {
|
||||
|
||||
this.clusterManager = clusterManager;
|
||||
try {
|
||||
listener.startListening();
|
||||
register();
|
||||
|
||||
} catch (IOException e) {
|
||||
throw new PersistenceException(e);
|
||||
}
|
||||
}
|
||||
|
||||
private void send(SocketClient client, SocketClusterMessage msg) {
|
||||
public void shutdown() {
|
||||
deregister();
|
||||
listener.shutdown();
|
||||
}
|
||||
|
||||
try {
|
||||
// alternative would be to connect/disconnect here
|
||||
// but prefer to use keepalive
|
||||
client.send(msg);
|
||||
|
||||
} catch (Exception ex){
|
||||
logger.error("Error sending message", ex);
|
||||
try {
|
||||
client.reconnect();
|
||||
} catch (IOException e) {
|
||||
logger.error("Error trying to reconnect", ex);
|
||||
}
|
||||
}
|
||||
/**
|
||||
* Register with all the other members of the Cluster.
|
||||
*/
|
||||
private void register() {
|
||||
|
||||
SocketClusterMessage h = SocketClusterMessage.register(local.getHostPort(), true);
|
||||
|
||||
for (int i = 0; i < members.length; i++) {
|
||||
boolean online = members[i].register(h);
|
||||
logger.info("Cluster Member [{}] online[{}]", members[i].getHostPort(), online);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Send the payload to all the members of the cluster.
|
||||
*/
|
||||
public void broadcast(RemoteTransactionEvent remoteTransEvent) {
|
||||
try {
|
||||
|
||||
txnOutgoing.incrementAndGet();
|
||||
DataHolder dataHolder = txnSerialiseHelper.createDataHolder(remoteTransEvent);
|
||||
SocketClusterMessage msg = SocketClusterMessage.transEvent(dataHolder);
|
||||
broadcast(msg);
|
||||
} catch (Exception e){
|
||||
String msg = "Error sending RemoteTransactionEvent "+remoteTransEvent+" to cluster members.";
|
||||
logger.error(msg, e);
|
||||
}
|
||||
protected void setMemberOnline(String fullName, boolean online) throws IOException {
|
||||
synchronized (clientMap) {
|
||||
logger.info("Cluster Member [{}] online[{}]", fullName, online);
|
||||
SocketClient member = clientMap.get(fullName);
|
||||
member.setOnline(online);
|
||||
}
|
||||
}
|
||||
|
||||
protected void broadcast(SocketClusterMessage msg) {
|
||||
|
||||
for (int i = 0; i < members.length; i++) {
|
||||
send(members[i], msg);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Leave the cluster.
|
||||
*/
|
||||
private void deregister() {
|
||||
|
||||
SocketClusterMessage h = SocketClusterMessage.register(local.getHostPort(), false);
|
||||
broadcast(h);
|
||||
|
||||
for (int i = 0; i < members.length; i++) {
|
||||
members[i].disconnect();
|
||||
}
|
||||
}
|
||||
private void send(SocketClient client, SocketClusterMessage msg) {
|
||||
|
||||
/**
|
||||
* Process a Cluster message.
|
||||
*/
|
||||
protected boolean process(SocketConnection request) throws IOException, ClassNotFoundException {
|
||||
try {
|
||||
// alternative would be to connect/disconnect here but prefer to use keepalive
|
||||
if (logger.isTraceEnabled()) {
|
||||
logger.trace("... send to member {} broadcast msg: {}", client, msg);
|
||||
}
|
||||
client.send(msg);
|
||||
|
||||
try {
|
||||
SocketClusterMessage h = (SocketClusterMessage)request.readObject();
|
||||
|
||||
if (h.isRegisterEvent()){
|
||||
setMemberOnline(h.getRegisterHost(), h.isRegister());
|
||||
|
||||
} else {
|
||||
txnIncoming.incrementAndGet();
|
||||
DataHolder dataHolder = h.getDataHolder();
|
||||
RemoteTransactionEvent transEvent = txnSerialiseHelper.read(dataHolder);
|
||||
transEvent.run();
|
||||
}
|
||||
|
||||
if (h.isRegisterEvent() && !h.isRegister()){
|
||||
// instance shutting down
|
||||
return true;
|
||||
} else {
|
||||
return false;
|
||||
}
|
||||
} catch (InterruptedIOException e) {
|
||||
String msg = "Timeout waiting for message";
|
||||
logger.info(msg, e);
|
||||
try {
|
||||
request.disconnect();
|
||||
} catch (IOException ex){
|
||||
logger.info("Error disconnecting after timeout", ex);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
} catch (Exception ex) {
|
||||
logger.error("Error sending message", ex);
|
||||
try {
|
||||
client.reconnect();
|
||||
} catch (IOException e) {
|
||||
logger.error("Error trying to reconnect", ex);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a host:port into a InetSocketAddress.
|
||||
*/
|
||||
private InetSocketAddress parseFullName(String hostAndPort) {
|
||||
|
||||
try {
|
||||
hostAndPort = hostAndPort.trim();
|
||||
int colonPos = hostAndPort.indexOf(":");
|
||||
if (colonPos == -1) {
|
||||
String msg = "No colon \":\" in "+hostAndPort;
|
||||
throw new IllegalArgumentException(msg);
|
||||
}
|
||||
String host = hostAndPort.substring(0, colonPos);
|
||||
String sPort = hostAndPort.substring(colonPos + 1, hostAndPort.length());
|
||||
int port = Integer.parseInt(sPort);
|
||||
|
||||
return new InetSocketAddress(host, port);
|
||||
|
||||
} catch (Exception ex){
|
||||
throw new RuntimeException("Error parsing ["+hostAndPort+"] for the form [host:port]", ex);
|
||||
}
|
||||
/**
|
||||
* Send the payload to all the members of the cluster.
|
||||
*/
|
||||
public void broadcast(RemoteTransactionEvent remoteTransEvent) {
|
||||
try {
|
||||
txnOutgoing.incrementAndGet();
|
||||
DataHolder dataHolder = txnSerialiseHelper.createDataHolder(remoteTransEvent);
|
||||
SocketClusterMessage msg = SocketClusterMessage.transEvent(dataHolder);
|
||||
broadcast(msg);
|
||||
} catch (Exception e) {
|
||||
logger.error("Error sending RemoteTransactionEvent " + remoteTransEvent + " to cluster members.", e);
|
||||
}
|
||||
|
||||
class TxnSerialiseHelper extends SerialiseTransactionHelper {
|
||||
}
|
||||
|
||||
@Override
|
||||
public SpiEbeanServer getEbeanServer(String serverName) {
|
||||
return (SpiEbeanServer)clusterManager.getServer(serverName);
|
||||
}
|
||||
protected void broadcast(SocketClusterMessage msg) {
|
||||
|
||||
if (logger.isTraceEnabled()) {
|
||||
logger.trace("... broadcast msg: "+msg);
|
||||
}
|
||||
for (int i = 0; i < members.length; i++) {
|
||||
send(members[i], msg);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Leave the cluster.
|
||||
*/
|
||||
private void deregister() {
|
||||
|
||||
SocketClusterMessage h = SocketClusterMessage.register(local.getHostPort(), false);
|
||||
broadcast(h);
|
||||
for (int i = 0; i < members.length; i++) {
|
||||
members[i].disconnect();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Process an incoming Cluster message.
|
||||
*/
|
||||
protected boolean process(SocketConnection request) throws IOException, ClassNotFoundException {
|
||||
|
||||
try {
|
||||
SocketClusterMessage h = (SocketClusterMessage) request.readObject();
|
||||
if (logger.isTraceEnabled()) {
|
||||
logger.trace("... received msg: {}", h);
|
||||
}
|
||||
|
||||
if (h.isRegisterEvent()) {
|
||||
setMemberOnline(h.getRegisterHost(), h.isRegister());
|
||||
|
||||
} else {
|
||||
txnIncoming.incrementAndGet();
|
||||
DataHolder dataHolder = h.getDataHolder();
|
||||
RemoteTransactionEvent transEvent = txnSerialiseHelper.read(dataHolder);
|
||||
transEvent.run();
|
||||
}
|
||||
|
||||
// instance shutting down
|
||||
return h.isRegisterEvent() && !h.isRegister();
|
||||
|
||||
} catch (InterruptedIOException e) {
|
||||
logger.info("Timeout waiting for message", e);
|
||||
try {
|
||||
request.disconnect();
|
||||
} catch (IOException ex) {
|
||||
logger.info("Error disconnecting after timeout", ex);
|
||||
}
|
||||
return true;
|
||||
|
||||
} catch (EOFException e) {
|
||||
logger.info("EOF disconnecting");
|
||||
return true;
|
||||
} catch (IOException e) {
|
||||
logger.info("IO Error waiting/reading message", e);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a host:port into a InetSocketAddress.
|
||||
*/
|
||||
private InetSocketAddress parseFullName(String hostAndPort) {
|
||||
|
||||
try {
|
||||
hostAndPort = hostAndPort.trim();
|
||||
int colonPos = hostAndPort.indexOf(":");
|
||||
if (colonPos == -1) {
|
||||
String msg = "No colon \":\" in " + hostAndPort;
|
||||
throw new IllegalArgumentException(msg);
|
||||
}
|
||||
String host = hostAndPort.substring(0, colonPos);
|
||||
String sPort = hostAndPort.substring(colonPos + 1, hostAndPort.length());
|
||||
int port = Integer.parseInt(sPort);
|
||||
|
||||
return new InetSocketAddress(host, port);
|
||||
|
||||
} catch (Exception ex) {
|
||||
throw new RuntimeException("Error parsing [" + hostAndPort + "] for the form [host:port]", ex);
|
||||
}
|
||||
}
|
||||
|
||||
class TxnSerialiseHelper extends SerialiseTransactionHelper {
|
||||
|
||||
@Override
|
||||
public SpiEbeanServer getEbeanServer(String serverName) {
|
||||
return (SpiEbeanServer) clusterManager.getServer(serverName);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+99
-122
@@ -6,10 +6,10 @@ import java.net.ServerSocket;
|
||||
import java.net.Socket;
|
||||
import java.net.SocketException;
|
||||
|
||||
import com.avaje.ebeaninternal.server.lib.DaemonThreadPool;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import com.avaje.ebeaninternal.server.lib.thread.ThreadPool;
|
||||
|
||||
/**
|
||||
* Serverside multithreaded socket listener. Accepts connections and dispatches
|
||||
@@ -26,141 +26,118 @@ import com.avaje.ebeaninternal.server.lib.thread.ThreadPool;
|
||||
*/
|
||||
class SocketClusterListener implements Runnable {
|
||||
|
||||
private static final Logger logger = LoggerFactory.getLogger(SocketClusterListener.class);
|
||||
|
||||
/**
|
||||
* The port the SocketListener uses.
|
||||
*/
|
||||
private final int port;
|
||||
private static final Logger logger = LoggerFactory.getLogger(SocketClusterListener.class);
|
||||
|
||||
/**
|
||||
* The length of the socket accept timeout.
|
||||
*/
|
||||
private final int listenTimeout = 60000;
|
||||
/**
|
||||
* The server socket used to listen for requests.
|
||||
*/
|
||||
private final ServerSocket serverListenSocket;
|
||||
|
||||
/**
|
||||
* The server socket used to listen for requests.
|
||||
*/
|
||||
private final ServerSocket serverListenSocket;
|
||||
/**
|
||||
* The listening thread.
|
||||
*/
|
||||
private final Thread listenerThread;
|
||||
|
||||
/**
|
||||
* The listening thread.
|
||||
*/
|
||||
private final Thread listenerThread;
|
||||
/**
|
||||
* The pool of threads that actually do the parsing execution of requests.
|
||||
*/
|
||||
private final DaemonThreadPool threadPool;
|
||||
|
||||
/**
|
||||
* The pool of threads that actually do the parsing execution of requests.
|
||||
*/
|
||||
private final ThreadPool threadPool;
|
||||
private final SocketClusterBroadcast owner;
|
||||
|
||||
private final SocketClusterBroadcast owner;
|
||||
/**
|
||||
* shutting down flag.
|
||||
*/
|
||||
boolean doingShutdown;
|
||||
|
||||
/**
|
||||
* shutting down flag.
|
||||
*/
|
||||
boolean doingShutdown;
|
||||
/**
|
||||
* Whether the listening thread is busy assigning a request to a thread.
|
||||
*/
|
||||
boolean isActive;
|
||||
|
||||
/**
|
||||
* Whether the listening thread is busy assigning a request to a thread.
|
||||
*/
|
||||
boolean isActive;
|
||||
/**
|
||||
* Construct with a given thread pool name.
|
||||
*/
|
||||
public SocketClusterListener(SocketClusterBroadcast owner, int port, int coreThreads, int maxThreads, String poolName) {
|
||||
this.owner = owner;
|
||||
this.threadPool = new DaemonThreadPool(coreThreads, maxThreads, 60, 30, poolName);
|
||||
try {
|
||||
this.serverListenSocket = new ServerSocket(port);
|
||||
this.serverListenSocket.setSoTimeout(60000);
|
||||
this.listenerThread = new Thread(this, "EbeanClusterListener");
|
||||
|
||||
/**
|
||||
* Construct with a given thread pool name.
|
||||
*/
|
||||
public SocketClusterListener(SocketClusterBroadcast owner, int port) {
|
||||
this.owner = owner;
|
||||
this.threadPool = ThreadPool.createThreadPool("EbeanCluster");
|
||||
this.port = port;
|
||||
|
||||
try {
|
||||
this.serverListenSocket = new ServerSocket(port);
|
||||
this.serverListenSocket.setSoTimeout(listenTimeout);
|
||||
this.listenerThread = new Thread(this, "EbeanClusterListener");
|
||||
|
||||
} catch (IOException e){
|
||||
String msg = "Error starting cluster socket listener on port "+port;
|
||||
throw new RuntimeException(msg,e);
|
||||
} catch (IOException e) {
|
||||
String msg = "Error starting cluster socket listener on port " + port;
|
||||
throw new RuntimeException(msg, e);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Start listening for requests.
|
||||
*/
|
||||
public void startListening() throws IOException {
|
||||
logger.trace("... startListening()");
|
||||
this.listenerThread.setDaemon(true);
|
||||
this.listenerThread.start();
|
||||
}
|
||||
|
||||
/**
|
||||
* Shutdown this listener.
|
||||
*/
|
||||
public void shutdown() {
|
||||
doingShutdown = true;
|
||||
try {
|
||||
if (isActive) {
|
||||
synchronized (listenerThread) {
|
||||
try {
|
||||
listenerThread.wait(1000);
|
||||
} catch (InterruptedException e) {
|
||||
// OK to ignore as expected to Interrupt for shutdown.
|
||||
}
|
||||
}
|
||||
}
|
||||
listenerThread.interrupt();
|
||||
serverListenSocket.close();
|
||||
} catch (IOException e) {
|
||||
logger.error("Error shutting down listener", e);
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the port the listener is using.
|
||||
*/
|
||||
public int getPort() {
|
||||
return port;
|
||||
}
|
||||
threadPool.shutdown();
|
||||
}
|
||||
|
||||
/**
|
||||
* Start listening for requests.
|
||||
*/
|
||||
public void startListening() throws IOException {
|
||||
this.listenerThread.setDaemon(true);
|
||||
this.listenerThread.start();
|
||||
}
|
||||
|
||||
/**
|
||||
* Shutdown this listener.
|
||||
*/
|
||||
public void shutdown() {
|
||||
doingShutdown = true;
|
||||
try {
|
||||
if (isActive) {
|
||||
synchronized (listenerThread) {
|
||||
try {
|
||||
listenerThread.wait(1000);
|
||||
} catch (InterruptedException e) {
|
||||
// OK to ignore as expected to Interrupt for shutdown.
|
||||
;
|
||||
}
|
||||
}
|
||||
}
|
||||
listenerThread.interrupt();
|
||||
serverListenSocket.close();
|
||||
} catch (IOException e) {
|
||||
logger.error("Error shutting down listener", e);
|
||||
/**
|
||||
* This is a runnable and so this must be public. Don't call this externally
|
||||
* but rather call the startListening() method.
|
||||
*/
|
||||
public void run() {
|
||||
// run in loop until doingShutdown is true...
|
||||
while (!doingShutdown) {
|
||||
try {
|
||||
synchronized (listenerThread) {
|
||||
Socket clientSocket = serverListenSocket.accept();
|
||||
isActive = true;
|
||||
|
||||
Runnable request = new RequestProcessor(owner, clientSocket);
|
||||
threadPool.execute(request);
|
||||
|
||||
isActive = false;
|
||||
}
|
||||
|
||||
threadPool.shutdown();
|
||||
}
|
||||
|
||||
/**
|
||||
* This is a runnable and so this must be public. Don't call this externally
|
||||
* but rather call the startListening() method.
|
||||
*/
|
||||
public void run() {
|
||||
// run in loop until doingShutdown is true...
|
||||
while (!doingShutdown) {
|
||||
try {
|
||||
synchronized (listenerThread) {
|
||||
Socket clientSocket = serverListenSocket.accept();
|
||||
|
||||
isActive = true;
|
||||
|
||||
Runnable request = new RequestProcessor(owner, clientSocket);
|
||||
threadPool.assign(request, true);
|
||||
|
||||
isActive = false;
|
||||
}
|
||||
} catch (SocketException e) {
|
||||
if (doingShutdown) {
|
||||
String msg = "doingShutdown and accept threw:"+ e.getMessage();
|
||||
logger.info(msg);
|
||||
|
||||
} else {
|
||||
logger.error(null, e);
|
||||
}
|
||||
|
||||
} catch (InterruptedIOException e) {
|
||||
// this will happen when the server is very quiet.
|
||||
// that is, no requests
|
||||
logger.debug("Possibly expected due to accept timeout?" + e.getMessage());
|
||||
|
||||
} catch (IOException e) {
|
||||
// log it and continue in the loop...
|
||||
logger.error(null, e);
|
||||
}
|
||||
} catch (SocketException e) {
|
||||
if (doingShutdown) {
|
||||
logger.info("doingShutdown and accept threw:" + e.getMessage());
|
||||
} else {
|
||||
logger.error("Error while listening", e);
|
||||
}
|
||||
} catch (InterruptedIOException e) {
|
||||
// this will happen when the server is very quiet.
|
||||
// that is, no requests
|
||||
logger.debug("Possibly expected due to accept timeout? {}", e.getMessage());
|
||||
|
||||
} catch (IOException e) {
|
||||
// log it and continue in the loop...
|
||||
logger.error("IOException processing cluster message", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user