mirror of
https://github.com/ebean-orm/ebean.git
synced 2024-04-21 10:51:47 +00:00
#1427 - QueryCache should be cleared, if one of a dependent bean is updated
Update RemoteTransactionEvent, add (back) common binary message handling (Common to TCP, K8s and upcoming binary pubsub)
This commit is contained in:
@@ -1,8 +1,8 @@
|
||||
package io.ebeaninternal.api;
|
||||
|
||||
import io.ebean.event.BulkTableEvent;
|
||||
import io.ebeaninternal.server.cluster.BinaryMessage;
|
||||
import io.ebeaninternal.server.cluster.BinaryMessageList;
|
||||
import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage;
|
||||
import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList;
|
||||
|
||||
import java.io.DataInput;
|
||||
import java.io.DataOutputStream;
|
||||
@@ -102,7 +102,7 @@ public final class TransactionEventTable implements Serializable {
|
||||
os.writeBoolean(insert);
|
||||
os.writeBoolean(update);
|
||||
os.writeBoolean(delete);
|
||||
|
||||
os.close();
|
||||
msgList.add(msg);
|
||||
}
|
||||
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
package io.ebeaninternal.server.cache;
|
||||
|
||||
import io.ebeaninternal.server.cluster.BinaryMessage;
|
||||
import io.ebeaninternal.server.cluster.BinaryMessageList;
|
||||
import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage;
|
||||
import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList;
|
||||
|
||||
import java.io.DataInput;
|
||||
import java.io.DataOutputStream;
|
||||
|
||||
@@ -13,7 +13,7 @@ import java.util.concurrent.ConcurrentHashMap;
|
||||
/**
|
||||
* Manages the cluster service.
|
||||
*/
|
||||
public class ClusterManager {
|
||||
public class ClusterManager implements MessageServerProvider {
|
||||
|
||||
private static final Logger clusterLogger = LoggerFactory.getLogger("io.ebean.Cluster");
|
||||
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
package io.ebeaninternal.server.cluster;
|
||||
|
||||
import io.ebean.EbeanServer;
|
||||
|
||||
/**
|
||||
* Returns EbeanServer instances for remote message reading.
|
||||
*/
|
||||
public interface MessageServerProvider {
|
||||
|
||||
/**
|
||||
* Return the EbeanServer instance by name.
|
||||
*/
|
||||
EbeanServer getServer(String name);
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
package io.ebeaninternal.server.cluster.binarymessage;
|
||||
|
||||
import io.ebeaninternal.api.SpiEbeanServer;
|
||||
import io.ebeaninternal.api.TransactionEventTable;
|
||||
import io.ebeaninternal.server.cache.RemoteCacheEvent;
|
||||
import io.ebeaninternal.server.cluster.MessageServerProvider;
|
||||
import io.ebeaninternal.server.transaction.BeanPersistIds;
|
||||
import io.ebeaninternal.server.transaction.RemoteTableMod;
|
||||
import io.ebeaninternal.server.transaction.RemoteTransactionEvent;
|
||||
|
||||
import java.io.ByteArrayInputStream;
|
||||
import java.io.DataInputStream;
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* Reads the binary message returning RemoteTransactionEvent.
|
||||
*/
|
||||
class BinaryDataReader {
|
||||
|
||||
private final MessageServerProvider serverProvider;
|
||||
private final DataInputStream dataInput;
|
||||
|
||||
private SpiEbeanServer server;
|
||||
private RemoteTransactionEvent event;
|
||||
|
||||
BinaryDataReader(MessageServerProvider serverProvider, byte[] data) {
|
||||
this.serverProvider = serverProvider;
|
||||
this.dataInput = new DataInputStream(new ByteArrayInputStream(data));
|
||||
}
|
||||
|
||||
/**
|
||||
* Read the binary message returning a RemoteTransactionEvent.
|
||||
*/
|
||||
RemoteTransactionEvent read() throws IOException {
|
||||
|
||||
String serverName = dataInput.readUTF();
|
||||
|
||||
this.server = (SpiEbeanServer) serverProvider.getServer(serverName);
|
||||
if (server == null) {
|
||||
throw new IllegalStateException("EbeanServer not found for name [" + serverName + "]");
|
||||
}
|
||||
this.event = new RemoteTransactionEvent(server);
|
||||
boolean more = dataInput.readBoolean();
|
||||
while (more) {
|
||||
readMessage();
|
||||
more = dataInput.readBoolean();
|
||||
}
|
||||
return event;
|
||||
}
|
||||
|
||||
private void readMessage() throws IOException {
|
||||
|
||||
int msgType = dataInput.readInt();
|
||||
switch (msgType) {
|
||||
case BinaryMessage.TYPE_BEANIUD:
|
||||
event.addBeanPersistIds(BeanPersistIds.readBinaryMessage(server, dataInput));
|
||||
break;
|
||||
|
||||
case BinaryMessage.TYPE_TABLEIUD:
|
||||
event.addTableIUD(TransactionEventTable.TableIUD.readBinaryMessage(dataInput));
|
||||
break;
|
||||
|
||||
case BinaryMessage.TYPE_CACHE:
|
||||
event.addRemoteCacheEvent(RemoteCacheEvent.readBinaryMessage(dataInput));
|
||||
break;
|
||||
|
||||
case BinaryMessage.TYPE_TABLEMOD:
|
||||
event.addRemoteTableMod(RemoteTableMod.readBinaryMessage(dataInput));
|
||||
break;
|
||||
|
||||
default:
|
||||
throw new RuntimeException("Invalid Transaction msgType " + msgType);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
package io.ebeaninternal.server.cluster.binarymessage;
|
||||
|
||||
import java.io.ByteArrayOutputStream;
|
||||
import java.io.DataOutputStream;
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* Writes BinaryMessageList to DataHolder raw byte[].
|
||||
*/
|
||||
class BinaryDataWriter {
|
||||
|
||||
private final BinaryMessageList messageList;
|
||||
|
||||
private final ByteArrayOutputStream buffer;
|
||||
|
||||
private final DataOutputStream dataOut;
|
||||
|
||||
private final String serverName;
|
||||
|
||||
BinaryDataWriter(String serverName, BinaryMessageList messageList) {
|
||||
this.serverName = serverName;
|
||||
this.messageList = messageList;
|
||||
this.buffer = new ByteArrayOutputStream(256);
|
||||
this.dataOut = new DataOutputStream(buffer);
|
||||
}
|
||||
|
||||
/**
|
||||
* Write the message as raw byte[].
|
||||
*/
|
||||
byte[] write() throws IOException {
|
||||
|
||||
//TODO: Add Protocol version and "podId" to support pubsub cluster
|
||||
//dataOut.writeInt(MsgKeys.MESSAGE_KEY);
|
||||
dataOut.writeUTF(serverName);
|
||||
for (BinaryMessage msg : messageList.getList()) {
|
||||
write(msg);
|
||||
}
|
||||
dataOut.writeBoolean(false);
|
||||
dataOut.flush();
|
||||
return buffer.toByteArray();
|
||||
}
|
||||
|
||||
private void write(BinaryMessage msg) throws IOException {
|
||||
|
||||
dataOut.writeBoolean(true);
|
||||
dataOut.write(msg.getByteArray());
|
||||
}
|
||||
|
||||
}
|
||||
+2
-1
@@ -1,4 +1,4 @@
|
||||
package io.ebeaninternal.server.cluster;
|
||||
package io.ebeaninternal.server.cluster.binarymessage;
|
||||
|
||||
import java.io.ByteArrayOutputStream;
|
||||
import java.io.DataOutputStream;
|
||||
@@ -24,6 +24,7 @@ public class BinaryMessage {
|
||||
public static final int TYPE_BEANIUD = 1;
|
||||
public static final int TYPE_TABLEIUD = 2;
|
||||
public static final int TYPE_CACHE = 3;
|
||||
public static final int TYPE_TABLEMOD = 4;
|
||||
|
||||
public static final int TYPE_MSGACK = 8;
|
||||
public static final int TYPE_MSGRESEND = 9;
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package io.ebeaninternal.server.cluster;
|
||||
package io.ebeaninternal.server.cluster.binarymessage;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
@@ -0,0 +1,147 @@
|
||||
package io.ebeaninternal.server.cluster.binarymessage;
|
||||
|
||||
import java.io.DataInputStream;
|
||||
import java.io.DataOutputStream;
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* The message broadcast around the cluster.
|
||||
*/
|
||||
public class ClusterMessage {
|
||||
|
||||
private static final int MAX_LENGTH = 10 * 1024 * 1024;
|
||||
|
||||
private final String registerIp;
|
||||
|
||||
private final String podName;
|
||||
|
||||
private final boolean register;
|
||||
|
||||
private final byte[] data;
|
||||
|
||||
/**
|
||||
* Create a register message.
|
||||
*/
|
||||
public static ClusterMessage register(String registerIp, boolean register, String podName) {
|
||||
return new ClusterMessage(registerIp, register, podName);
|
||||
}
|
||||
|
||||
/**
|
||||
* Create a transaction message.
|
||||
*/
|
||||
public static ClusterMessage transEvent(byte[] data) {
|
||||
return new ClusterMessage(data);
|
||||
}
|
||||
|
||||
/**
|
||||
* Create for register online/offline message.
|
||||
*/
|
||||
private ClusterMessage(String registerIp, boolean register, String podName) {
|
||||
this.registerIp = registerIp;
|
||||
this.register = register;
|
||||
this.podName = podName;
|
||||
this.data = null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Create for a transaction message.
|
||||
*/
|
||||
private ClusterMessage(byte[] data) {
|
||||
this.data = data;
|
||||
this.registerIp = null;
|
||||
this.podName = null;
|
||||
this.register = false;
|
||||
}
|
||||
|
||||
public String toString() {
|
||||
StringBuilder sb = new StringBuilder();
|
||||
if (registerIp != null) {
|
||||
sb.append("register ");
|
||||
sb.append(register);
|
||||
sb.append(" ");
|
||||
sb.append(registerIp);
|
||||
} else {
|
||||
sb.append("[data]");
|
||||
}
|
||||
return sb.toString();
|
||||
}
|
||||
|
||||
/**
|
||||
* Return true if this is a register event as opposed to a transaction message.
|
||||
*/
|
||||
public boolean isRegisterEvent() {
|
||||
return registerIp != null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Return the register host for online/offline message.
|
||||
*/
|
||||
public String getRegisterIp() {
|
||||
return registerIp;
|
||||
}
|
||||
|
||||
public String getPodName() {
|
||||
return podName;
|
||||
}
|
||||
|
||||
/**
|
||||
* Return true if register is true for a online/offline message.
|
||||
*/
|
||||
public boolean isRegister() {
|
||||
return register;
|
||||
}
|
||||
|
||||
/**
|
||||
* Return the raw message data.
|
||||
*/
|
||||
public byte[] getData() {
|
||||
return data;
|
||||
}
|
||||
|
||||
/**
|
||||
* Write the message in binary form.
|
||||
*/
|
||||
public void write(DataOutputStream dataOutput) throws IOException {
|
||||
|
||||
if (data != null) {
|
||||
// write data message
|
||||
dataOutput.writeInt(MsgKeys.DATA);
|
||||
dataOutput.writeInt(data.length);
|
||||
dataOutput.write(data);
|
||||
} else {
|
||||
// write header message
|
||||
dataOutput.writeInt(MsgKeys.HEADER);
|
||||
dataOutput.writeUTF(getRegisterIp());
|
||||
dataOutput.writeBoolean(register);
|
||||
dataOutput.writeUTF(getPodName());
|
||||
}
|
||||
dataOutput.flush();
|
||||
}
|
||||
|
||||
/**
|
||||
* Read the message from binary form.
|
||||
*/
|
||||
public static ClusterMessage read(DataInputStream dataInput) throws IOException, InvalidMessageException {
|
||||
|
||||
int key = dataInput.readInt();
|
||||
if (key == MsgKeys.DATA) {
|
||||
int length = dataInput.readInt();
|
||||
if (length > MAX_LENGTH) {
|
||||
throw new IOException("Message data too large length:"+length);
|
||||
}
|
||||
byte[] data = new byte[length];
|
||||
dataInput.readFully(data);
|
||||
return new ClusterMessage(data);
|
||||
|
||||
} else if (key == MsgKeys.HEADER) {
|
||||
String host = dataInput.readUTF();
|
||||
boolean registered = dataInput.readBoolean();
|
||||
String podName = dataInput.readUTF();
|
||||
return new ClusterMessage(host, registered, podName);
|
||||
|
||||
} else {
|
||||
throw new InvalidMessageException("Invalid message key:" + key);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
+8
@@ -0,0 +1,8 @@
|
||||
package io.ebeaninternal.server.cluster.binarymessage;
|
||||
|
||||
public class InvalidMessageException extends Exception {
|
||||
|
||||
public InvalidMessageException(String message) {
|
||||
super(message);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
package io.ebeaninternal.server.cluster.binarymessage;
|
||||
|
||||
import io.ebeaninternal.server.cluster.MessageServerProvider;
|
||||
import io.ebeaninternal.server.transaction.RemoteTransactionEvent;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* Mechanism to convert RemoteTransactionEvent to/from byte[] content.
|
||||
*/
|
||||
public class MessageReadWrite {
|
||||
|
||||
private final MessageServerProvider serverProvider;
|
||||
|
||||
public MessageReadWrite(MessageServerProvider serverProvider) {
|
||||
this.serverProvider = serverProvider;
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert the RemoteTransactionEvent to raw byte[] content.
|
||||
*/
|
||||
public byte[] write(RemoteTransactionEvent transEvent) throws IOException {
|
||||
|
||||
BinaryMessageList messageList = new BinaryMessageList();
|
||||
transEvent.writeBinaryMessage(messageList);
|
||||
|
||||
return new BinaryDataWriter(transEvent.getServerName(), messageList).write();
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert the byte[] content to RemoteTransactionEvent.
|
||||
*/
|
||||
public RemoteTransactionEvent read(byte[] data) throws IOException {
|
||||
|
||||
return new BinaryDataReader(serverProvider, data).read();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
package io.ebeaninternal.server.cluster.binarymessage;
|
||||
|
||||
public interface MsgKeys {
|
||||
|
||||
/**
|
||||
* Used to identify client on connection initiation.
|
||||
*/
|
||||
int HELLO = 182;
|
||||
|
||||
/**
|
||||
* Used to confirm protocol on reading messages.
|
||||
*/
|
||||
int HEADER = 11;
|
||||
|
||||
/**
|
||||
* Used to confirm protocol on reading messages.
|
||||
*/
|
||||
int DATA = 12;
|
||||
}
|
||||
@@ -1,8 +1,8 @@
|
||||
package io.ebeaninternal.server.transaction;
|
||||
|
||||
import io.ebeaninternal.api.SpiEbeanServer;
|
||||
import io.ebeaninternal.server.cluster.BinaryMessage;
|
||||
import io.ebeaninternal.server.cluster.BinaryMessageList;
|
||||
import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage;
|
||||
import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList;
|
||||
import io.ebeaninternal.server.core.PersistRequest;
|
||||
import io.ebeaninternal.server.deploy.BeanDescriptor;
|
||||
import io.ebeaninternal.server.deploy.id.IdBinder;
|
||||
@@ -140,7 +140,7 @@ public class BeanPersistIds {
|
||||
idBinder.writeData(os, idList.get(i));
|
||||
}
|
||||
|
||||
os.flush();
|
||||
os.close();
|
||||
msgList.add(m);
|
||||
|
||||
} while (i < eof);
|
||||
@@ -167,7 +167,7 @@ public class BeanPersistIds {
|
||||
return sb.toString();
|
||||
}
|
||||
|
||||
void addId(PersistRequest.Type type, Serializable id) {
|
||||
public void addId(PersistRequest.Type type, Serializable id) {
|
||||
switch (type) {
|
||||
case INSERT:
|
||||
addInsertId(id);
|
||||
@@ -210,10 +210,18 @@ public class BeanPersistIds {
|
||||
return beanDescriptor;
|
||||
}
|
||||
|
||||
List<Object> getDeleteIds() {
|
||||
public List<Object> getDeleteIds() {
|
||||
return deleteIds;
|
||||
}
|
||||
|
||||
public List<Object> getInsertIds() {
|
||||
return insertIds;
|
||||
}
|
||||
|
||||
public List<Object> getUpdateIds() {
|
||||
return updateIds;
|
||||
}
|
||||
|
||||
/**
|
||||
* Notify the cache of this event that came from another server in the cluster.
|
||||
*/
|
||||
|
||||
@@ -211,6 +211,10 @@ final class PostCommitProcessing {
|
||||
|
||||
RemoteTransactionEvent remoteTransactionEvent = new RemoteTransactionEvent(serverName);
|
||||
|
||||
Set<String> touched = cacheChanges.touchedTables();
|
||||
if (touched != null && !touched.isEmpty()) {
|
||||
remoteTransactionEvent.addRemoteTableMod(new RemoteTableMod(cacheChanges.modificationTimestamp(), touched));
|
||||
}
|
||||
if (beanPersistIdMap != null) {
|
||||
for (BeanPersistIds beanPersist : beanPersistIdMap.values()) {
|
||||
remoteTransactionEvent.addBeanPersistIds(beanPersist);
|
||||
|
||||
@@ -0,0 +1,55 @@
|
||||
package io.ebeaninternal.server.transaction;
|
||||
|
||||
import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage;
|
||||
import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList;
|
||||
|
||||
import java.io.DataInput;
|
||||
import java.io.DataOutputStream;
|
||||
import java.io.IOException;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.Set;
|
||||
|
||||
public class RemoteTableMod {
|
||||
|
||||
private final long timestamp;
|
||||
|
||||
private final Set<String> tables;
|
||||
|
||||
public RemoteTableMod(long timestamp, Set<String> tables) {
|
||||
this.timestamp = timestamp;
|
||||
this.tables = tables;
|
||||
}
|
||||
|
||||
public long getTimestamp() {
|
||||
return timestamp;
|
||||
}
|
||||
|
||||
public Set<String> getTables() {
|
||||
return tables;
|
||||
}
|
||||
|
||||
public static RemoteTableMod readBinaryMessage(DataInput dataInput) throws IOException {
|
||||
|
||||
long timestamp = dataInput.readLong();
|
||||
int count = dataInput.readInt();
|
||||
|
||||
Set<String> tables = new LinkedHashSet<>();
|
||||
for (int i = 0; i < count; i++) {
|
||||
tables.add(dataInput.readUTF());
|
||||
}
|
||||
return new RemoteTableMod(timestamp, tables);
|
||||
}
|
||||
|
||||
public void writeBinary(BinaryMessageList msgList) throws IOException {
|
||||
BinaryMessage msg = new BinaryMessage(tables.size() * 20 + 30);
|
||||
DataOutputStream os = msg.getOs();
|
||||
os.writeInt(BinaryMessage.TYPE_TABLEMOD);
|
||||
os.writeLong(timestamp);
|
||||
os.writeInt(tables.size());
|
||||
for (String table : tables) {
|
||||
os.writeUTF(table);
|
||||
}
|
||||
os.close();
|
||||
msgList.add(msg);
|
||||
}
|
||||
}
|
||||
@@ -3,7 +3,7 @@ package io.ebeaninternal.server.transaction;
|
||||
import io.ebeaninternal.api.SpiEbeanServer;
|
||||
import io.ebeaninternal.api.TransactionEventTable.TableIUD;
|
||||
import io.ebeaninternal.server.cache.RemoteCacheEvent;
|
||||
import io.ebeaninternal.server.cluster.BinaryMessageList;
|
||||
import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.ArrayList;
|
||||
@@ -19,6 +19,8 @@ public class RemoteTransactionEvent implements Runnable {
|
||||
|
||||
private RemoteCacheEvent remoteCacheEvent;
|
||||
|
||||
private RemoteTableMod remoteTableMod;
|
||||
|
||||
private String serverName;
|
||||
|
||||
private transient SpiEbeanServer server;
|
||||
@@ -53,6 +55,10 @@ public class RemoteTransactionEvent implements Runnable {
|
||||
|
||||
public void writeBinaryMessage(BinaryMessageList msgList) throws IOException {
|
||||
|
||||
if (remoteTableMod != null) {
|
||||
remoteTableMod.writeBinary(msgList);
|
||||
}
|
||||
|
||||
if (tableList != null) {
|
||||
for (TableIUD aTableList : tableList) {
|
||||
aTableList.writeBinaryMessage(msgList);
|
||||
@@ -113,6 +119,10 @@ public class RemoteTransactionEvent implements Runnable {
|
||||
tableList.add(tableIud);
|
||||
}
|
||||
|
||||
public void addRemoteTableMod(RemoteTableMod remoteTableMod) {
|
||||
this.remoteTableMod = remoteTableMod;
|
||||
}
|
||||
|
||||
public String getServerName() {
|
||||
return serverName;
|
||||
}
|
||||
@@ -141,4 +151,8 @@ public class RemoteTransactionEvent implements Runnable {
|
||||
return remoteCacheEvent;
|
||||
}
|
||||
|
||||
public RemoteTableMod getRemoteTableMod() {
|
||||
return remoteTableMod;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user