diff --git a/pom.xml b/pom.xml
index 04ffadca0..47c5d9d19 100644
--- a/pom.xml
+++ b/pom.xml
@@ -9,7 +9,7 @@
io.ebean
ebean
- 11.6.2-SNAPSHOT
+ 11.7.1-SNAPSHOT
jar
ebean
diff --git a/src/main/java/io/ebean/CacheMode.java b/src/main/java/io/ebean/CacheMode.java
index 364f06e3f..f8c6cfe70 100644
--- a/src/main/java/io/ebean/CacheMode.java
+++ b/src/main/java/io/ebean/CacheMode.java
@@ -56,17 +56,29 @@ public enum CacheMode {
AUTO(true, true),
/**
- * Do not read from cache, but load retrieved beans into the cache.
+ * Do not read from cache, but put beans into the cache and invalidate parts of the cache as necessary.
*
- * Use this, if you want to get the fresh value from database into the cache. Typically a CacheMode.ON query
- * will follow.
+ * Use this on a query if you want to get the fresh value from database and put it into the cache.
*/
+ PUT(false, true),
+
+ /**
+ * Deprecated - migrate to PUT.
+ */
+ @Deprecated
RECACHE(false, true),
/**
- * Query the cache for value. If it is there, use it and otherwise hit database but do NOT put the value
- * into the cache. Note that there are not many use case for this mode.
+ * GET only from the cache.
+ *
+ * This mode does not put entries into the cache or invalidate parts of the cache.
*/
+ GET(true, false),
+
+ /**
+ * Deprecated - migrate to GET.
+ */
+ @Deprecated
QUERY_ONLY(true, false);
private boolean get;
diff --git a/src/main/java/io/ebean/cache/ServerCacheManager.java b/src/main/java/io/ebean/cache/ServerCacheManager.java
index 1ee97f7a3..38de72fc6 100644
--- a/src/main/java/io/ebean/cache/ServerCacheManager.java
+++ b/src/main/java/io/ebean/cache/ServerCacheManager.java
@@ -44,4 +44,19 @@ public interface ServerCacheManager {
*/
void clearAll();
+ /**
+ * Clear all the local caches.
+ *
+ * This is used when the L2 Cache is based on clustered near-caches (Like Ebean-K8s-L2Cache).
+ * It is not used when the L2 cache is a distributed cache such as HazelCast or Ignite etc.
+ */
+ void clearAllLocal();
+
+ /**
+ * Clear the local caches for this bean type.
+ *
+ * This is used when the L2 Cache is based on clustered near-caches (Like Ebean-K8s-L2Cache).
+ * It is not used when the L2 cache is a distributed cache such as HazelCast or Ignite etc.
+ */
+ void clearLocal(Class> beanType);
}
diff --git a/src/main/java/io/ebean/config/ContainerConfig.java b/src/main/java/io/ebean/config/ContainerConfig.java
index faaa49df8..32fb98770 100644
--- a/src/main/java/io/ebean/config/ContainerConfig.java
+++ b/src/main/java/io/ebean/config/ContainerConfig.java
@@ -10,22 +10,82 @@ import java.util.Properties;
*/
public class ContainerConfig {
- protected boolean clusterActive;
+ private boolean active;
+ private String serviceName;
+ private String namespace;
+ private String podName;
+ private int port;
- protected Properties properties;
+ private Properties properties;
+
+ /**
+ * Return the service name.
+ */
+ public String getServiceName() {
+ return serviceName;
+ }
+
+ /**
+ * Set the service name.
+ */
+ public void setServiceName(String serviceName) {
+ this.serviceName = serviceName;
+ }
+
+ /**
+ * Return the namespace.
+ */
+ public String getNamespace() {
+ return namespace;
+ }
+
+ /**
+ * Set the namespace.
+ */
+ public void setNamespace(String namespace) {
+ this.namespace = namespace;
+ }
+
+ /**
+ * Return the pod name.
+ */
+ public String getPodName() {
+ return podName;
+ }
+
+ /**
+ * Set the pod name.
+ */
+ public void setPodName(String podName) {
+ this.podName = podName;
+ }
+
+ /**
+ * Return the port to use.
+ */
+ public int getPort() {
+ return port;
+ }
+
+ /**
+ * Set the port to use.
+ */
+ public void setPort(int port) {
+ this.port = port;
+ }
/**
* Return true if clustering is active.
*/
- public boolean isClusterActive() {
- return clusterActive;
+ public boolean isActive() {
+ return active;
}
/**
* Set to true for clustering to be active.
*/
- public void setClusterActive(boolean clusterActive) {
- this.clusterActive = clusterActive;
+ public void setActive(boolean active) {
+ this.active = active;
}
/**
@@ -47,7 +107,14 @@ public class ContainerConfig {
*/
public void loadFromProperties(Properties properties) {
this.properties = properties;
- this.clusterActive = getProperty(properties, "ebean.cluster.active", clusterActive);
+ this.active = getProperty(properties, "ebean.cluster.active", active);
+ this.serviceName = properties.getProperty("ebean.cluster.serviceName", serviceName);
+ this.namespace = properties.getProperty("ebean.cluster.namespace", namespace);
+ this.podName = properties.getProperty("ebean.cluster.podName", podName);
+ String portParam = properties.getProperty("ebean.cluster.port");
+ if (portParam != null) {
+ this.port = Integer.parseInt(portParam);
+ }
}
/**
diff --git a/src/main/java/io/ebeaninternal/api/NaturalKeyEntry.java b/src/main/java/io/ebeaninternal/api/NaturalKeyEntry.java
index 50a096d7a..756b24006 100644
--- a/src/main/java/io/ebeaninternal/api/NaturalKeyEntry.java
+++ b/src/main/java/io/ebeaninternal/api/NaturalKeyEntry.java
@@ -49,8 +49,10 @@ public class NaturalKeyEntry {
}
private void load(List eqList) {
- for (NaturalKeyEq eq : eqList) {
- map.put(eq.property, eq.value);
+ if (eqList != null) {
+ for (NaturalKeyEq eq : eqList) {
+ map.put(eq.property, eq.value);
+ }
}
}
diff --git a/src/main/java/io/ebeaninternal/server/cache/CacheManagerOptions.java b/src/main/java/io/ebeaninternal/server/cache/CacheManagerOptions.java
new file mode 100644
index 000000000..b61b48080
--- /dev/null
+++ b/src/main/java/io/ebeaninternal/server/cache/CacheManagerOptions.java
@@ -0,0 +1,85 @@
+package io.ebeaninternal.server.cache;
+
+import io.ebean.cache.ServerCacheFactory;
+import io.ebean.cache.ServerCacheOptions;
+import io.ebean.config.CurrentTenantProvider;
+import io.ebean.config.ServerConfig;
+import io.ebeaninternal.server.cluster.ClusterManager;
+
+/**
+ * Configuration options when creating the default cache manager.
+ */
+public class CacheManagerOptions {
+
+ private final ClusterManager clusterManager;
+
+ private final ServerConfig serverConfig;
+
+ private final boolean localL2Caching;
+
+ private CurrentTenantProvider currentTenantProvider;
+
+ private ServerCacheFactory cacheFactory = new DefaultServerCacheFactory();
+ private ServerCacheOptions beanDefault = new ServerCacheOptions();
+ private ServerCacheOptions queryDefault = new ServerCacheOptions();
+
+ CacheManagerOptions() {
+ this.localL2Caching = true;
+ this.clusterManager = null;
+ this.serverConfig = null;
+ this.cacheFactory = new DefaultServerCacheFactory();
+ this.beanDefault = new ServerCacheOptions();
+ this.queryDefault = new ServerCacheOptions();
+ }
+
+ public CacheManagerOptions(ClusterManager clusterManager, ServerConfig serverConfig, boolean localL2Caching) {
+ this.clusterManager = clusterManager;
+ this.serverConfig = serverConfig;
+ this.localL2Caching = localL2Caching;
+ this.currentTenantProvider = serverConfig.getCurrentTenantProvider();
+ }
+
+ public CacheManagerOptions with(ServerCacheOptions beanDefault, ServerCacheOptions queryDefault) {
+ this.beanDefault = beanDefault;
+ this.queryDefault = queryDefault;
+ return this;
+ }
+
+ public CacheManagerOptions with(ServerCacheFactory cacheFactory) {
+ this.cacheFactory = cacheFactory;
+ return this;
+ }
+
+ public CacheManagerOptions with(CurrentTenantProvider currentTenantProvider) {
+ this.currentTenantProvider = currentTenantProvider;
+ return this;
+ }
+
+ public String getServerName() {
+ return (serverConfig == null) ? "db" : serverConfig.getName();
+ }
+
+ public boolean isLocalL2Caching() {
+ return localL2Caching;
+ }
+
+ public ServerCacheFactory getCacheFactory() {
+ return cacheFactory;
+ }
+
+ public ServerCacheOptions getBeanDefault() {
+ return beanDefault;
+ }
+
+ public ServerCacheOptions getQueryDefault() {
+ return queryDefault;
+ }
+
+ public CurrentTenantProvider getCurrentTenantProvider() {
+ return currentTenantProvider;
+ }
+
+ public ClusterManager getClusterManager() {
+ return clusterManager;
+ }
+}
diff --git a/src/main/java/io/ebeaninternal/server/cache/DefaultCacheAdapter.java b/src/main/java/io/ebeaninternal/server/cache/DefaultCacheAdapter.java
index b5c47b054..0c3cfa9e5 100644
--- a/src/main/java/io/ebeaninternal/server/cache/DefaultCacheAdapter.java
+++ b/src/main/java/io/ebeaninternal/server/cache/DefaultCacheAdapter.java
@@ -44,12 +44,21 @@ public class DefaultCacheAdapter implements ServerCacheManager {
@Override
public void clear(Class> beanType) {
- cacheManager.getBeanCache(beanType).clear();
- cacheManager.getQueryCache(beanType).clear();
+ cacheManager.clear(beanType);
}
@Override
public void clearAll() {
cacheManager.clearAll();
}
+
+ @Override
+ public void clearAllLocal() {
+ cacheManager.clearAllLocal();
+ }
+
+ @Override
+ public void clearLocal(Class> beanType) {
+ cacheManager.clearLocal(beanType);
+ }
}
diff --git a/src/main/java/io/ebeaninternal/server/cache/DefaultCacheHolder.java b/src/main/java/io/ebeaninternal/server/cache/DefaultCacheHolder.java
index 1d4f43b74..092fe7b9a 100644
--- a/src/main/java/io/ebeaninternal/server/cache/DefaultCacheHolder.java
+++ b/src/main/java/io/ebeaninternal/server/cache/DefaultCacheHolder.java
@@ -8,16 +8,24 @@ import io.ebean.cache.ServerCacheOptions;
import io.ebean.cache.ServerCacheType;
import io.ebean.config.CurrentTenantProvider;
import io.ebean.util.AnnotationUtil;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentSkipListSet;
/**
* Manages the construction of caches.
*/
class DefaultCacheHolder {
+ private static final Logger log = LoggerFactory.getLogger("io.ebean.cache.ALL");
+
private final ConcurrentHashMap allCaches = new ConcurrentHashMap<>();
+ private final ConcurrentHashMap> collectIdCaches = new ConcurrentHashMap<>();
+
private final ServerCacheFactory cacheFactory;
private final ServerCacheOptions beanDefault;
@@ -25,6 +33,10 @@ class DefaultCacheHolder {
private final CurrentTenantProvider tenantProvider;
+ DefaultCacheHolder(CacheManagerOptions builder) {
+ this(builder.getCacheFactory(), builder.getBeanDefault(), builder.getQueryDefault(), builder.getCurrentTenantProvider());
+ }
+
/**
* Create with a cache factory and default cache options.
*
@@ -42,7 +54,6 @@ class DefaultCacheHolder {
ServerCache getCache(Class> beanType, String cacheKey, ServerCacheType type) {
return getCacheInternal(beanType, cacheKey, type);
-
}
private String key(String cacheKey, ServerCacheType type) {
@@ -60,16 +71,43 @@ class DefaultCacheHolder {
private ServerCache createCache(Class> beanType, ServerCacheType type, String key) {
ServerCacheOptions options = getCacheOptions(beanType, type);
-
+ if (type == ServerCacheType.COLLECTION_IDS) {
+ synchronized (this) {
+ collectIdCaches.computeIfAbsent(beanType.getName(), s -> new ConcurrentSkipListSet<>()).add(key);
+ }
+ }
return cacheFactory.createCache(type, key, tenantProvider, options);
}
void clearAll() {
+ log.debug("clearAll");
for (ServerCache serverCache : allCaches.values()) {
serverCache.clear();
}
}
+
+ public void clear(String name) {
+ log.debug("clear {}", name);
+ clearIfExists(key(name, ServerCacheType.QUERY));
+ clearIfExists(key(name, ServerCacheType.BEAN));
+ clearIfExists(key(name, ServerCacheType.NATURAL_KEY));
+ Set keys = collectIdCaches.get(name);
+ if (keys != null) {
+ for (String collectionIdKey : keys) {
+ clearIfExists(collectionIdKey);
+ }
+ }
+ }
+
+ private void clearIfExists(String fullKey) {
+ ServerCache cache = allCaches.get(fullKey);
+ if (cache != null) {
+ log.trace("clear cache {}", fullKey);
+ cache.clear();
+ }
+ }
+
/**
* Return the cache options for a given bean type.
*/
diff --git a/src/main/java/io/ebeaninternal/server/cache/DefaultServerCacheManager.java b/src/main/java/io/ebeaninternal/server/cache/DefaultServerCacheManager.java
index c8b52a6f7..bdb10d453 100644
--- a/src/main/java/io/ebeaninternal/server/cache/DefaultServerCacheManager.java
+++ b/src/main/java/io/ebeaninternal/server/cache/DefaultServerCacheManager.java
@@ -1,34 +1,37 @@
package io.ebeaninternal.server.cache;
import io.ebean.cache.ServerCache;
-import io.ebean.cache.ServerCacheFactory;
-import io.ebean.cache.ServerCacheOptions;
import io.ebean.cache.ServerCacheType;
-import io.ebean.config.CurrentTenantProvider;
+import io.ebeaninternal.server.cluster.ClusterManager;
/**
* Manages the bean and query caches.
*/
public class DefaultServerCacheManager implements SpiCacheManager {
+ private final ClusterManager clusterManager;
+
private final DefaultCacheHolder cacheHolder;
private final boolean localL2Caching;
+ private final String serverName;
+
/**
* Create with a cache factory and default cache options.
*/
- public DefaultServerCacheManager(boolean localL2Caching, CurrentTenantProvider tenantProvider, ServerCacheFactory cacheFactory,
- ServerCacheOptions beanDefault, ServerCacheOptions queryDefault) {
- this.localL2Caching = localL2Caching;
- this.cacheHolder = new DefaultCacheHolder(cacheFactory, beanDefault, queryDefault, tenantProvider);
+ public DefaultServerCacheManager(CacheManagerOptions builder) {
+ this.clusterManager = builder.getClusterManager();
+ this.serverName = builder.getServerName();
+ this.localL2Caching = builder.isLocalL2Caching();
+ this.cacheHolder = new DefaultCacheHolder(builder);
}
/**
* Construct when l2 cache is disabled.
*/
public DefaultServerCacheManager() {
- this(true, null, new DefaultServerCacheFactory(), new ServerCacheOptions(), new ServerCacheOptions());
+ this(new CacheManagerOptions());
}
@Override
@@ -42,6 +45,27 @@ public class DefaultServerCacheManager implements SpiCacheManager {
@Override
public void clearAll() {
cacheHolder.clearAll();
+ if (clusterManager != null) {
+ clusterManager.cacheClearAll(serverName);
+ }
+ }
+
+ @Override
+ public void clearAllLocal() {
+ cacheHolder.clearAll();
+ }
+
+ @Override
+ public void clear(Class> beanType) {
+ cacheHolder.clear(name(beanType));
+ if (clusterManager != null) {
+ clusterManager.cacheClear(serverName, beanType);
+ }
+ }
+
+ @Override
+ public void clearLocal(Class> beanType) {
+ cacheHolder.clear(name(beanType));
}
@Override
diff --git a/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java b/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java
new file mode 100644
index 000000000..4881cfd77
--- /dev/null
+++ b/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java
@@ -0,0 +1,93 @@
+package io.ebeaninternal.server.cache;
+
+import io.ebeaninternal.server.cluster.BinaryMessage;
+import io.ebeaninternal.server.cluster.BinaryMessageList;
+
+import java.io.DataInput;
+import java.io.DataOutputStream;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.List;
+
+/**
+ * Cache events broadcast across the cluster.
+ */
+public class RemoteCacheEvent {
+
+ private boolean clearAll;
+
+ private List clearCaches;
+
+ /**
+ * Clear all the caches.
+ */
+ public RemoteCacheEvent(boolean clearAll) {
+ this.clearAll = clearAll;
+ this.clearCaches = null;
+ }
+
+ /**
+ * Clear caches for the given bean type.
+ */
+ public RemoteCacheEvent(Class> beanType) {
+ this.clearAll = false;
+ this.clearCaches = new ArrayList<>(1);
+ this.clearCaches.add(beanType.getName());
+ }
+
+ /**
+ * Create with the read options.
+ */
+ public RemoteCacheEvent(boolean clearAll, List beanTypes) {
+ this.clearAll = clearAll;
+ this.clearCaches = beanTypes;
+ }
+
+ public boolean isClearAll() {
+ return clearAll;
+ }
+
+ public List getClearCaches() {
+ return clearCaches;
+ }
+
+ public String toString() {
+ return "clearAll:" + clearAll + " caches:" + clearCaches;
+ }
+
+ public static RemoteCacheEvent readBinaryMessage(DataInput dataInput) throws IOException {
+
+ boolean clearAll = dataInput.readBoolean();
+ int size = dataInput.readInt();
+
+ List clearCache = null;
+ if (size > 0) {
+ clearCache = new ArrayList<>(size);
+ for (int i = 0; i < size; i++) {
+ clearCache.add(dataInput.readUTF());
+ }
+ }
+
+ return new RemoteCacheEvent(clearAll, clearCache);
+ }
+
+ public void writeBinaryMessage(BinaryMessageList msgList) throws IOException {
+
+ int bufferSize = (clearCaches == null) ? 0 : clearCaches.size() * 30;
+
+ BinaryMessage msg = new BinaryMessage(bufferSize + 10);
+ DataOutputStream os = msg.getOs();
+ os.writeInt(BinaryMessage.TYPE_CACHE);
+ os.writeBoolean(clearAll);
+ if (clearCaches == null) {
+ os.writeInt(0);
+ } else {
+ os.writeInt(clearCaches.size());
+ for (String cacheName : clearCaches) {
+ os.writeUTF(cacheName);
+ }
+ }
+ msgList.add(msg);
+
+ }
+}
diff --git a/src/main/java/io/ebeaninternal/server/cache/SpiCacheManager.java b/src/main/java/io/ebeaninternal/server/cache/SpiCacheManager.java
index bd80836e7..3e83203bd 100644
--- a/src/main/java/io/ebeaninternal/server/cache/SpiCacheManager.java
+++ b/src/main/java/io/ebeaninternal/server/cache/SpiCacheManager.java
@@ -36,9 +36,24 @@ public interface SpiCacheManager {
*/
ServerCache getQueryCache(Class> beanType);
+ /**
+ * Clear the caches for the given bean type.
+ */
+ void clear(Class> beanType);
+
/**
* Clear all the caches.
*/
void clearAll();
+ /**
+ * Clear all local caches.
+ */
+ void clearAllLocal();
+
+ /**
+ * Clear local caches for the given bean type.
+ */
+ void clearLocal(Class> beanType);
+
}
diff --git a/src/main/java/io/ebeaninternal/server/cluster/BinaryMessage.java b/src/main/java/io/ebeaninternal/server/cluster/BinaryMessage.java
index c470c1793..b78b3f7eb 100644
--- a/src/main/java/io/ebeaninternal/server/cluster/BinaryMessage.java
+++ b/src/main/java/io/ebeaninternal/server/cluster/BinaryMessage.java
@@ -23,6 +23,7 @@ public class BinaryMessage {
public static final int TYPE_MSGCONTROL = 0;
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_MSGACK = 8;
public static final int TYPE_MSGRESEND = 9;
diff --git a/src/main/java/io/ebeaninternal/server/cluster/BinaryMessageList.java b/src/main/java/io/ebeaninternal/server/cluster/BinaryMessageList.java
index 3b3513063..f84b7995c 100644
--- a/src/main/java/io/ebeaninternal/server/cluster/BinaryMessageList.java
+++ b/src/main/java/io/ebeaninternal/server/cluster/BinaryMessageList.java
@@ -5,12 +5,10 @@ import java.util.List;
/**
* Holds a List of BinaryMessage's.
- *
- * @author rbygrave
*/
public class BinaryMessageList {
- final ArrayList list = new ArrayList<>();
+ final List list = new ArrayList<>();
public void add(BinaryMessage msg) {
list.add(msg);
diff --git a/src/main/java/io/ebeaninternal/server/cluster/ClusterBroadcast.java b/src/main/java/io/ebeaninternal/server/cluster/ClusterBroadcast.java
index bef3d3aef..fbfd153f0 100644
--- a/src/main/java/io/ebeaninternal/server/cluster/ClusterBroadcast.java
+++ b/src/main/java/io/ebeaninternal/server/cluster/ClusterBroadcast.java
@@ -24,5 +24,4 @@ public interface ClusterBroadcast {
* Send a transaction event to all the members of the cluster.
*/
void broadcast(RemoteTransactionEvent remoteTransEvent);
-
}
diff --git a/src/main/java/io/ebeaninternal/server/cluster/ClusterBroadcastFactory.java b/src/main/java/io/ebeaninternal/server/cluster/ClusterBroadcastFactory.java
index b350e036f..07391c4fc 100644
--- a/src/main/java/io/ebeaninternal/server/cluster/ClusterBroadcastFactory.java
+++ b/src/main/java/io/ebeaninternal/server/cluster/ClusterBroadcastFactory.java
@@ -1,6 +1,6 @@
package io.ebeaninternal.server.cluster;
-import java.util.Properties;
+import io.ebean.config.ContainerConfig;
/**
* Factory to create the cluster broadcast service.
@@ -10,5 +10,5 @@ public interface ClusterBroadcastFactory {
/**
* Create the cluster transport with the manager and deployment properties.
*/
- ClusterBroadcast create(ClusterManager manager, Properties properties);
+ ClusterBroadcast create(ClusterManager manager, ContainerConfig config);
}
diff --git a/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java b/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java
index aa8f2fcf4..6fd56c713 100644
--- a/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java
+++ b/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java
@@ -17,8 +17,6 @@ public class ClusterManager {
private static final Logger clusterLogger = LoggerFactory.getLogger("io.ebean.Cluster");
- private static final Logger logger = LoggerFactory.getLogger(ClusterManager.class);
-
private final ConcurrentHashMap serverMap = new ConcurrentHashMap<>();
private final Object monitor = new Object();
@@ -27,12 +25,14 @@ public class ClusterManager {
private boolean started;
+ private boolean shutdown;
+
public ClusterManager(ContainerConfig config) {
- if (!config.isClusterActive()) {
- broadcast = null;
+ ClusterBroadcastFactory factory = createFactory();
+ if (factory != null) {
+ broadcast = factory.create(this, config);
} else {
- ClusterBroadcastFactory factory = createFactory();
- broadcast = factory.create(this, config.getProperties());
+ broadcast = null;
}
}
@@ -47,10 +47,6 @@ public class ClusterManager {
if (iterator.hasNext()) {
factory = iterator.next();
}
- if (factory == null) {
- throw new IllegalStateException("No ClusterTransportFactory found in classpath. "
- + " Probably need to add the avaje-ebeanorm-cluster dependency");
- }
return factory;
}
@@ -76,6 +72,24 @@ public class ClusterManager {
}
}
+ /**
+ * Broadcast a cache clear all event to the cluster.
+ */
+ public void cacheClearAll(String serverName) {
+ if (broadcast != null) {
+ broadcast.broadcast(new RemoteTransactionEvent(serverName).cacheClearAll());
+ }
+ }
+
+ /**
+ * Broadcast a cache clear event to the cluster.
+ */
+ public void cacheClear(String serverName, Class> beanType) {
+ if (broadcast != null) {
+ broadcast.broadcast(new RemoteTransactionEvent(serverName).cacheClear(beanType));
+ }
+ }
+
/**
* Return true if clustering is on.
*/
@@ -99,8 +113,8 @@ public class ClusterManager {
* Shutdown the service and Deregister from the cluster.
*/
public void shutdown() {
- if (broadcast != null) {
- logger.info("ClusterManager shutdown ");
+ if (broadcast != null && !shutdown) {
+ shutdown = true;
broadcast.shutdown();
}
}
diff --git a/src/main/java/io/ebeaninternal/server/core/DefaultContainer.java b/src/main/java/io/ebeaninternal/server/core/DefaultContainer.java
index 6566b5554..3afb745ea 100644
--- a/src/main/java/io/ebeaninternal/server/core/DefaultContainer.java
+++ b/src/main/java/io/ebeaninternal/server/core/DefaultContainer.java
@@ -4,7 +4,6 @@ import io.ebean.BackgroundExecutor;
import io.ebean.cache.ServerCacheFactory;
import io.ebean.cache.ServerCacheOptions;
import io.ebean.cache.ServerCachePlugin;
-import io.ebean.service.SpiContainer;
import io.ebean.config.ContainerConfig;
import io.ebean.config.PropertyMap;
import io.ebean.config.ServerConfig;
@@ -12,10 +11,12 @@ import io.ebean.config.TenantMode;
import io.ebean.config.UnderscoreNamingConvention;
import io.ebean.config.dbplatform.DatabasePlatform;
import io.ebean.config.dbplatform.h2.H2Platform;
-import io.ebeaninternal.dbmigration.DbOffline;
+import io.ebean.service.SpiContainer;
import io.ebeaninternal.api.SpiBackgroundExecutor;
import io.ebeaninternal.api.SpiContainerBootup;
import io.ebeaninternal.api.SpiEbeanServer;
+import io.ebeaninternal.dbmigration.DbOffline;
+import io.ebeaninternal.server.cache.CacheManagerOptions;
import io.ebeaninternal.server.cache.DefaultServerCacheManager;
import io.ebeaninternal.server.cache.DefaultServerCachePlugin;
import io.ebeaninternal.server.cache.SpiCacheManager;
@@ -71,6 +72,7 @@ public class DefaultContainer implements SpiContainer {
@Override
public void shutdown() {
+ clusterManager.shutdown();
ShutdownManager.shutdown();
}
@@ -201,7 +203,11 @@ public class DefaultContainer implements SpiContainer {
}
ServerCacheFactory factory = plugin.create(serverConfig, executor);
- return new DefaultServerCacheManager(localL2Caching, serverConfig.getCurrentTenantProvider(), factory, beanOptions, queryOptions);
+
+ CacheManagerOptions builder = new CacheManagerOptions(clusterManager, serverConfig, localL2Caching)
+ .with(beanOptions, queryOptions)
+ .with(factory);
+ return new DefaultServerCacheManager(builder);
}
/**
diff --git a/src/main/java/io/ebeaninternal/server/core/DefaultServer.java b/src/main/java/io/ebeaninternal/server/core/DefaultServer.java
index b5192582b..64cb8cec7 100644
--- a/src/main/java/io/ebeaninternal/server/core/DefaultServer.java
+++ b/src/main/java/io/ebeaninternal/server/core/DefaultServer.java
@@ -65,6 +65,7 @@ import io.ebeaninternal.api.TransactionEventTable;
import io.ebeaninternal.dbmigration.DdlGenerator;
import io.ebeaninternal.dbmigration.ddlgeneration.DdlHandler;
import io.ebeaninternal.server.autotune.AutoTuneService;
+import io.ebeaninternal.server.cache.RemoteCacheEvent;
import io.ebeaninternal.server.core.timezone.DataTimeZone;
import io.ebeaninternal.server.deploy.BeanDescriptor;
import io.ebeaninternal.server.deploy.BeanDescriptorManager;
@@ -2072,6 +2073,30 @@ public final class DefaultServer implements SpiServer, SpiEbeanServer {
@Override
public void remoteTransactionEvent(RemoteTransactionEvent event) {
transactionManager.remoteTransactionEvent(event);
+ processRemoteCacheEvent(event);
+ }
+
+ /**
+ * Process a cache event coming from another server in the cluster.
+ */
+ private void processRemoteCacheEvent(RemoteTransactionEvent event) {
+ RemoteCacheEvent cacheEvent = event.getRemoteCacheEvent();
+ if (cacheEvent != null) {
+ if (cacheEvent.isClearAll()) {
+ serverCacheManager.clearAllLocal();
+ } else {
+ List caches = cacheEvent.getClearCaches();
+ if (caches != null) {
+ for (String cache : caches) {
+ try {
+ serverCacheManager.clearLocal(Class.forName(cache));
+ } catch (Exception e) {
+ logger.error("Error clearing local cache for type " + cache, e);
+ }
+ }
+ }
+ }
+ }
}
private P executeInTrans(Function fun, Transaction t) {
diff --git a/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java b/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java
index 6a4ffc630..f6047d209 100644
--- a/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java
+++ b/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java
@@ -2,6 +2,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 java.io.IOException;
@@ -16,6 +17,8 @@ public class RemoteTransactionEvent implements Runnable {
private DeleteByIdMap deleteByIdMap;
+ private RemoteCacheEvent remoteCacheEvent;
+
private String serverName;
private transient SpiEbeanServer server;
@@ -65,6 +68,9 @@ public class RemoteTransactionEvent implements Runnable {
for (BeanPersistIds aBeanPersistList : beanPersistList) {
aBeanPersistList.writeBinaryMessage(msgList);
}
+ if (remoteCacheEvent != null) {
+ remoteCacheEvent.writeBinaryMessage(msgList);
+ }
}
public boolean isEmpty() {
@@ -77,6 +83,29 @@ public class RemoteTransactionEvent implements Runnable {
beanPersistList.add(beanPersist);
}
+ /**
+ * Add a cache clearAll event.
+ */
+ public RemoteTransactionEvent cacheClearAll() {
+ this.remoteCacheEvent = new RemoteCacheEvent(true);
+ return this;
+ }
+
+ /**
+ * Add a cache clear event for the given bean type.
+ */
+ public RemoteTransactionEvent cacheClear(Class> beanType) {
+ this.remoteCacheEvent = new RemoteCacheEvent(beanType);
+ return this;
+ }
+
+ /**
+ * Set the RemoteCacheEvent.
+ */
+ public void addRemoteCacheEvent(RemoteCacheEvent remoteCacheEvent) {
+ this.remoteCacheEvent = remoteCacheEvent;
+ }
+
public void addTableIUD(TableIUD tableIud) {
if (tableList == null) {
tableList = new ArrayList<>(4);
@@ -108,4 +137,8 @@ public class RemoteTransactionEvent implements Runnable {
return beanPersistList;
}
+ public RemoteCacheEvent getRemoteCacheEvent() {
+ return remoteCacheEvent;
+ }
+
}
diff --git a/src/main/java/io/ebeaninternal/server/type/ArrayElementConverterEnum.java b/src/main/java/io/ebeaninternal/server/type/ArrayElementConverterEnum.java
new file mode 100644
index 000000000..12f3aa624
--- /dev/null
+++ b/src/main/java/io/ebeaninternal/server/type/ArrayElementConverterEnum.java
@@ -0,0 +1,23 @@
+package io.ebeaninternal.server.type;
+
+public class ArrayElementConverterEnum implements ArrayElementConverter {
+
+ final ScalarType> scalarType;
+
+ final Class extends Enum>> valueType1;
+
+ public ArrayElementConverterEnum(ScalarType> scalarType, Class extends Enum>> valueType1) {
+ this.scalarType = scalarType;
+ this.valueType1 = valueType1;
+ }
+
+ @Override
+ public String toElement(Object rawValue) {
+
+ Enum>[] enumConstants = valueType1.getEnumConstants();
+ if (scalarType == null) {
+ return rawValue.toString();
+ }
+ return scalarType.format(rawValue);
+ }
+}
diff --git a/src/test/java/io/ebean/BaseTestCase.java b/src/test/java/io/ebean/BaseTestCase.java
index 7f98b2941..effac9205 100644
--- a/src/test/java/io/ebean/BaseTestCase.java
+++ b/src/test/java/io/ebean/BaseTestCase.java
@@ -129,7 +129,7 @@ public abstract class BaseTestCase {
protected void loadCountryCache() {
Ebean.find(Country.class)
- .setBeanCacheMode(CacheMode.RECACHE)
+ .setBeanCacheMode(CacheMode.PUT)
.findList();
}
diff --git a/src/test/java/io/ebean/config/ContainerConfigTest.java b/src/test/java/io/ebean/config/ContainerConfigTest.java
new file mode 100644
index 000000000..23c4d0b6c
--- /dev/null
+++ b/src/test/java/io/ebean/config/ContainerConfigTest.java
@@ -0,0 +1,31 @@
+package io.ebean.config;
+
+import org.junit.Test;
+
+import java.util.Properties;
+
+import static org.junit.Assert.assertEquals;
+
+public class ContainerConfigTest {
+
+ @Test
+ public void loadFromProperties() {
+
+ Properties p = new Properties();
+ p.setProperty("ebean.cluster.active", "true");
+ p.setProperty("ebean.cluster.serviceName", "a");
+ p.setProperty("ebean.cluster.namespace", "b");
+ p.setProperty("ebean.cluster.podName", "c");
+ p.setProperty("ebean.cluster.port", "42");
+
+
+ ContainerConfig containerConfig = new ContainerConfig();
+ containerConfig.loadFromProperties(p);
+
+ assertEquals(true, containerConfig.isActive());
+ assertEquals("a", containerConfig.getServiceName());
+ assertEquals("b", containerConfig.getNamespace());
+ assertEquals("c", containerConfig.getPodName());
+ assertEquals(42, containerConfig.getPort());
+ }
+}
diff --git a/src/test/java/io/ebeaninternal/server/cache/DefaultServerCacheManagerTest.java b/src/test/java/io/ebeaninternal/server/cache/DefaultServerCacheManagerTest.java
index 5abefa2a2..2cf1b9cbb 100644
--- a/src/test/java/io/ebeaninternal/server/cache/DefaultServerCacheManagerTest.java
+++ b/src/test/java/io/ebeaninternal/server/cache/DefaultServerCacheManagerTest.java
@@ -1,10 +1,12 @@
package io.ebeaninternal.server.cache;
-import io.ebean.cache.ServerCacheFactory;
-import io.ebean.cache.ServerCacheOptions;
+import io.ebean.config.ContainerConfig;
+import io.ebean.config.CurrentTenantProvider;
+import io.ebean.config.ServerConfig;
+import io.ebeaninternal.server.cluster.ClusterManager;
+import org.junit.Test;
import org.tests.model.basic.Contact;
import org.tests.model.basic.Customer;
-import org.junit.Test;
import static org.assertj.core.api.StrictAssertions.assertThat;
import static org.junit.Assert.assertTrue;
@@ -13,15 +15,31 @@ public class DefaultServerCacheManagerTest {
private ThreadLocal tenantId = new ThreadLocal<>();
- private final ServerCacheFactory cacheFactory = new DefaultServerCacheFactory();
+ class TdTenPro implements CurrentTenantProvider {
- private DefaultServerCacheManager manager = new DefaultServerCacheManager(true, null, cacheFactory, new ServerCacheOptions(), new ServerCacheOptions());
+ @Override
+ public Object currentId() {
+ return tenantId.get();
+ }
+ }
+
+ private ClusterManager clusterManager = new ClusterManager(new ContainerConfig());
+
+ private DefaultServerCacheManager manager = new DefaultServerCacheManager(new CacheManagerOptions(clusterManager, new ServerConfig(), true));
+
+ private DefaultServerCacheManager multiTenantManager;
+
+ public DefaultServerCacheManagerTest(){
+
+ CacheManagerOptions builder = new CacheManagerOptions(clusterManager, new ServerConfig(), true);
+ builder.with(new TdTenPro());
+
+ this.multiTenantManager = new DefaultServerCacheManager(builder);
+ }
- private DefaultServerCacheManager multiTenantManager = new DefaultServerCacheManager(true, tenantId::get, cacheFactory, new ServerCacheOptions(), new ServerCacheOptions());
@Test
- public void getCache_normal() throws Exception {
-
+ public void getCache_normal() {
DefaultServerCache cache = cache(manager, Customer.class);
assertThat(cache.getName()).isEqualTo("org.tests.model.basic.Customer_B");
@@ -42,6 +60,18 @@ public class DefaultServerCacheManagerTest {
DefaultServerCache collCache = (DefaultServerCache) manager.getCollectionIdsCache(Customer.class, "contacts");
assertThat(collCache.getName()).isEqualTo("org.tests.model.basic.Customer.contacts_C");
+
+ cache.clearCount.sumThenReset();
+ collCache.clearCount.sumThenReset();
+ queryCache.clearCount.sumThenReset();
+ natKeyCache.clearCount.sumThenReset();
+
+ manager.clear(Customer.class);
+
+ assertThat(cache.clearCount.sumThenReset()).isEqualTo(1);
+ assertThat(natKeyCache.clearCount.sumThenReset()).isEqualTo(1);
+ assertThat(queryCache.clearCount.sumThenReset()).isEqualTo(1);
+ assertThat(collCache.clearCount.sumThenReset()).isEqualTo(1);
}
private DefaultServerCache cache(DefaultServerCacheManager manager, Class> beanType) {
@@ -49,7 +79,7 @@ public class DefaultServerCacheManagerTest {
}
@Test
- public void getCache_multiTenant() throws Exception {
+ public void getCache_multiTenant() {
tenantId.set("ten1");
DefaultServerCache cache = cache(multiTenantManager, Customer.class);
@@ -57,32 +87,32 @@ public class DefaultServerCacheManagerTest {
cache.put("1", "tenant1");
tenantId.set("ten2");
-
+
assertThat(cache.get("1")).isNull();
tenantId.set("ten1");
-
+
assertThat(cache.get("1")).isNotNull();
}
@Test
- public void getCache_singleTenant() throws Exception {
+ public void getCache_singleTenant() {
tenantId.set("ten1");
DefaultServerCache cache = cache(manager, Customer.class);
assertThat(cache.getName()).isEqualTo("org.tests.model.basic.Customer_B");
-
+
cache.put("1", "tenant1");
tenantId.set("ten2");
-
+
assertThat(cache.get("1")).isEqualTo("tenant1");
-
+
}
@Test
- public void isLocalL2Caching() throws Exception {
+ public void isLocalL2Caching() {
assertTrue(manager.isLocalL2Caching());
assertTrue(multiTenantManager.isLocalL2Caching());