diff --git a/src/main/java/io/ebeaninternal/api/TransactionEvent.java b/src/main/java/io/ebeaninternal/api/TransactionEvent.java index 711fa49dd..d8b25c113 100644 --- a/src/main/java/io/ebeaninternal/api/TransactionEvent.java +++ b/src/main/java/io/ebeaninternal/api/TransactionEvent.java @@ -3,7 +3,9 @@ package io.ebeaninternal.api; import io.ebeaninternal.server.cache.CacheChangeSet; import io.ebeaninternal.server.core.PersistRequestBean; import io.ebeaninternal.server.deploy.BeanDescriptor; +import io.ebeaninternal.server.deploy.BeanDescriptorManager; import io.ebeaninternal.server.transaction.DeleteByIdMap; +import io.ebeaninternal.server.transaction.TransactionManager; import io.ebeanservice.docstore.api.DocStoreUpdates; import java.io.Serializable; @@ -27,8 +29,6 @@ public class TransactionEvent implements Serializable { */ private final transient boolean local; - private final long modificationTimestamp; - private TransactionEventTable eventTables; private transient TransactionEventBeans eventBeans; @@ -38,9 +38,8 @@ public class TransactionEvent implements Serializable { /** * Create the TransactionEvent, one per Transaction. */ - public TransactionEvent(long modificationTimestamp) { + public TransactionEvent() { this.local = true; - this.modificationTimestamp = modificationTimestamp; } public void addDeleteById(BeanDescriptor desc, Object id) { @@ -111,8 +110,20 @@ public class TransactionEvent implements Serializable { /** * Build and return the cache changeSet. */ - public CacheChangeSet buildCacheChanges() { - CacheChangeSet changeSet = new CacheChangeSet(modificationTimestamp); + public CacheChangeSet buildCacheChanges(TransactionManager manager) { + + if (eventBeans == null && deleteByIdMap == null && eventTables == null) { + return null; + } + + CacheChangeSet changeSet = new CacheChangeSet(manager.clockNowMillis()); + if (eventTables != null && !eventTables.isEmpty()) { + // notify cache with table based changes + BeanDescriptorManager dm = manager.getBeanDescriptorManager(); + for (TransactionEventTable.TableIUD tableIUD : eventTables.values()) { + dm.cacheNotify(tableIUD, changeSet); + } + } if (eventBeans != null) { eventBeans.notifyCache(changeSet); } diff --git a/src/main/java/io/ebeaninternal/server/cache/CacheChangeBeanUpdate.java b/src/main/java/io/ebeaninternal/server/cache/CacheChangeBeanUpdate.java index cce2a4f9b..5463eab5e 100644 --- a/src/main/java/io/ebeaninternal/server/cache/CacheChangeBeanUpdate.java +++ b/src/main/java/io/ebeaninternal/server/cache/CacheChangeBeanUpdate.java @@ -25,6 +25,6 @@ class CacheChangeBeanUpdate implements CacheChange { @Override public void apply() { - desc.cacheBeanUpdate(id, changes, updateNaturalKey, version); + desc.cacheApplyBeanUpdate(id, changes, updateNaturalKey, version); } } diff --git a/src/main/java/io/ebeaninternal/server/cache/CacheChangeSet.java b/src/main/java/io/ebeaninternal/server/cache/CacheChangeSet.java index d917756d1..5865d2bbc 100644 --- a/src/main/java/io/ebeaninternal/server/cache/CacheChangeSet.java +++ b/src/main/java/io/ebeaninternal/server/cache/CacheChangeSet.java @@ -22,6 +22,8 @@ public class CacheChangeSet { private final Set> queryCaches = new HashSet<>(); + private final Set> beanCaches = new HashSet<>(); + private final Map, CacheChangeBeanRemove> beanRemoveMap = new HashMap<>(); private final Map manyChangeMap = new HashMap<>(); @@ -51,6 +53,9 @@ public class CacheChangeSet { for (BeanDescriptor entry : queryCaches) { entry.clearQueryCache(); } + for (BeanDescriptor entry : beanCaches) { + entry.clearBeanCache(); + } for (CacheChange entry : entries) { entry.apply(); } @@ -69,6 +74,13 @@ public class CacheChangeSet { touchedTables.add(descriptor.getBaseTable()); } + /** + * Add invalidation on a set of tables. + */ + public void addInvalidate(Set tables) { + touchedTables.addAll(tables); + } + /** * Add an entry to clear a query cache. */ @@ -76,6 +88,13 @@ public class CacheChangeSet { queryCaches.add(descriptor); } + /** + * Add an entry to clear a bean cache. + */ + public void addClearBean(BeanDescriptor descriptor) { + beanCaches.add(descriptor); + } + /** * Add many property clear. */ diff --git a/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java b/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java index 264af7bc8..8c9420911 100644 --- a/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java +++ b/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java @@ -384,8 +384,8 @@ public final class OrmQueryRequest extends BeanRequest implements SpiOrmQuery } private int notifyCache(int rows, boolean update) { - if (rows > 0 && beanDescriptor.isCaching()) { - transaction.getEvent().add(beanDescriptor.getBaseTable(), false, update, !update); + if (rows > 0) { + beanDescriptor.cacheUpdateQuery(update, transaction); } return rows; } diff --git a/src/main/java/io/ebeaninternal/server/core/PersistRequestBean.java b/src/main/java/io/ebeaninternal/server/core/PersistRequestBean.java index 3c658f34d..4839790ec 100644 --- a/src/main/java/io/ebeaninternal/server/core/PersistRequestBean.java +++ b/src/main/java/io/ebeaninternal/server/core/PersistRequestBean.java @@ -450,14 +450,14 @@ public final class PersistRequestBean extends PersistRequest implements BeanP if (notifyCache) { switch (type) { case INSERT: - beanDescriptor.cacheHandleInsert(this, changeSet); + beanDescriptor.cachePersistInsert(this, changeSet); break; case UPDATE: - beanDescriptor.cacheHandleUpdate(idValue, this, changeSet); + beanDescriptor.cachePersistUpdate(idValue, this, changeSet); break; case DELETE: case DELETE_SOFT: - beanDescriptor.cacheHandleDelete(idValue, this, changeSet); + beanDescriptor.cachePersistDelete(idValue, this, changeSet); break; default: throw new IllegalStateException("Invalid type " + type); diff --git a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptor.java b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptor.java index 96285d3f9..3edb14519 100644 --- a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptor.java +++ b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptor.java @@ -40,6 +40,7 @@ import io.ebeaninternal.api.ConcurrencyMode; import io.ebeaninternal.api.LoadContext; import io.ebeaninternal.api.SpiEbeanServer; import io.ebeaninternal.api.SpiQuery; +import io.ebeaninternal.api.SpiTransaction; import io.ebeaninternal.api.SpiUpdatePlan; import io.ebeaninternal.api.TransactionEventTable.TableIUD; import io.ebeaninternal.api.json.SpiJsonReader; @@ -828,8 +829,11 @@ public class BeanDescriptor implements BeanType, STreeType { */ @SuppressWarnings("unchecked") public void initialiseDocMapping() { - for (BeanPropertyAssocMany aPropertiesMany : propertiesMany) { - aPropertiesMany.initialisePostTarget(); + for (BeanPropertyAssocMany many : propertiesMany) { + many.initialisePostTarget(); + } + for (BeanPropertyAssocOne one : propertiesOne) { + one.initialisePostTarget(); } if (inheritInfo != null && !inheritInfo.isRoot()) { docStoreAdapter = (DocStoreBeanAdapter) inheritInfo.getRoot().desc().docStoreAdapter(); @@ -1366,13 +1370,6 @@ public class BeanDescriptor implements BeanType, STreeType { cacheHelp.queryCachePut(id, entry); } - /** - * Add a query cache clear into the changeSet. - */ - public void queryCacheClear(CacheChangeSet changeSet) { - cacheHelp.queryCacheClear(changeSet); - } - /** * Try to load the beanCollection from cache return true if successful. */ @@ -1469,13 +1466,6 @@ public class BeanDescriptor implements BeanType, STreeType { return cacheHelp.beanCacheGet(id, readOnly, context); } - /** - * Remove a bean from the cache given its Id. - */ - public void cacheHandleDeleteByIds(Collection ids, CacheChangeSet changeSet) { - cacheHelp.handleDeleteIds(ids, changeSet); - } - /** * Remove a collection of beans from the cache given the ids. */ @@ -1510,38 +1500,52 @@ public class BeanDescriptor implements BeanType, STreeType { cacheHelp.cacheNaturalKeyPut(id, newKey); } + /** + * Check if bulk update or delete query has a cache impact. + */ + public void cacheUpdateQuery(boolean update, SpiTransaction transaction) { + cacheHelp.cacheUpdateQuery(update, transaction); + } + /** * Invalidate parts of cache due to SqlUpdate or external modification etc. */ - public void cacheHandleBulkUpdate(TableIUD tableIUD) { - cacheHelp.handleBulkUpdate(tableIUD); + public void cachePersistTableIUD(TableIUD tableIUD, CacheChangeSet changeSet) { + cacheHelp.persistTableIUD(tableIUD, changeSet); } /** * Remove a bean from the cache given its Id. */ - public void cacheHandleDelete(Object id, PersistRequestBean deleteRequest, CacheChangeSet changeSet) { - cacheHelp.handleDelete(id, deleteRequest, changeSet); + public void cachePersistDeleteByIds(Collection ids, CacheChangeSet changeSet) { + cacheHelp.persistDeleteIds(ids, changeSet); + } + + /** + * Remove a bean from the cache given its Id. + */ + public void cachePersistDelete(Object id, PersistRequestBean deleteRequest, CacheChangeSet changeSet) { + cacheHelp.persistDelete(id, deleteRequest, changeSet); } /** * Add the insert changes to the changeSet. */ - public void cacheHandleInsert(PersistRequestBean insertRequest, CacheChangeSet changeSet) { - cacheHelp.handleInsert(insertRequest, changeSet); + public void cachePersistInsert(PersistRequestBean insertRequest, CacheChangeSet changeSet) { + cacheHelp.persistInsert(insertRequest, changeSet); } /** * Add the update to the changeSet. */ - public void cacheHandleUpdate(Object id, PersistRequestBean updateRequest, CacheChangeSet changeSet) { - cacheHelp.handleUpdate(id, updateRequest, changeSet); + public void cachePersistUpdate(Object id, PersistRequestBean updateRequest, CacheChangeSet changeSet) { + cacheHelp.persistUpdate(id, updateRequest, changeSet); } /** * Apply the update to the cache. */ - public void cacheBeanUpdate(Object id, Map changes, boolean updateNaturalKey, long version) { + public void cacheApplyBeanUpdate(Object id, Map changes, boolean updateNaturalKey, long version) { cacheHelp.cacheBeanUpdate(id, changes, updateNaturalKey, version); } diff --git a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorCacheHelp.java b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorCacheHelp.java index 79cdf45e7..3066e14b4 100644 --- a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorCacheHelp.java +++ b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorCacheHelp.java @@ -7,6 +7,7 @@ import io.ebean.bean.PersistenceContext; import io.ebean.cache.QueryCacheEntry; import io.ebean.cache.ServerCache; import io.ebeaninternal.api.BeanCacheResult; +import io.ebeaninternal.api.SpiTransaction; import io.ebeaninternal.api.TransactionEventTable.TableIUD; import io.ebeaninternal.server.cache.CacheChangeSet; import io.ebeaninternal.server.cache.CachedBeanData; @@ -68,6 +69,8 @@ final class BeanDescriptorCacheHelp { private final ServerCache naturalKeyCache; private final ServerCache queryCache; + private final boolean noCaching; + /** * Set to true if all persist changes need to notify the cache. */ @@ -108,6 +111,7 @@ final class BeanDescriptorCacheHelp { this.beanCache = null; this.naturalKeyCache = null; } + this.noCaching = (beanCache == null && queryCache == null); } /** @@ -132,7 +136,7 @@ final class BeanDescriptorCacheHelp { */ private boolean isNotifyOnDeletes() { for (BeanPropertyAssocOne imported : propertiesOneImported) { - if (imported.isCacheNotify()) { + if (imported.isCacheNotifyRelationship()) { return true; } } @@ -722,10 +726,16 @@ final class BeanDescriptorCacheHelp { return true; } + void cacheUpdateQuery(boolean update, SpiTransaction transaction) { + if (invalidateQueryCache || cacheNotifyOnAll || (!update && cacheNotifyOnDelete)) { + transaction.getEvent().add(desc.getBaseTable(), false, update, !update); + } + } + /** * Add appropriate cache changes to support delete by id. */ - void handleDeleteIds(Collection ids, CacheChangeSet changeSet) { + void persistDeleteIds(Collection ids, CacheChangeSet changeSet) { if (invalidateQueryCache) { changeSet.addInvalidate(desc); } else { @@ -739,7 +749,7 @@ final class BeanDescriptorCacheHelp { /** * Add appropriate cache changes to support delete bean. */ - void handleDelete(Object id, PersistRequestBean deleteRequest, CacheChangeSet changeSet) { + void persistDelete(Object id, PersistRequestBean deleteRequest, CacheChangeSet changeSet) { if (invalidateQueryCache) { changeSet.addInvalidate(desc); } else { @@ -754,7 +764,7 @@ final class BeanDescriptorCacheHelp { /** * Add appropriate cache changes to support insert. */ - void handleInsert(PersistRequestBean insertRequest, CacheChangeSet changeSet) { + void persistInsert(PersistRequestBean insertRequest, CacheChangeSet changeSet) { if (invalidateQueryCache) { changeSet.addInvalidate(desc); } else { @@ -773,7 +783,7 @@ final class BeanDescriptorCacheHelp { /** * Add appropriate changes to support update. */ - void handleUpdate(Object id, PersistRequestBean updateRequest, CacheChangeSet changeSet) { + void persistUpdate(Object id, PersistRequestBean updateRequest, CacheChangeSet changeSet) { if (invalidateQueryCache) { changeSet.addInvalidate(desc); @@ -803,16 +813,24 @@ final class BeanDescriptorCacheHelp { /** * Invalidate parts of cache due to SqlUpdate or external modification etc. */ - void handleBulkUpdate(TableIUD tableIUD) { + void persistTableIUD(TableIUD tableIUD, CacheChangeSet changeSet) { + if (invalidateQueryCache) { + changeSet.addInvalidate(desc); + return; + } + if (noCaching) { + return; + } + changeSet.addInvalidate(desc); // inserts don't invalidate the bean cache if (tableIUD.isUpdateOrDelete()) { - beanCacheClear(); + changeSet.addClearBean(desc); } // any change invalidates the query cache - queryCacheClear(); + changeSet.addClearQuery(desc); // any change invalidates the collection IDs cache for (BeanPropertyAssocOne imported : propertiesOneImported) { - imported.cacheClear(); + imported.cacheClear(changeSet); } } diff --git a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorManager.java b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorManager.java index a517dd25c..470df5ac8 100644 --- a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorManager.java +++ b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorManager.java @@ -25,6 +25,7 @@ import io.ebean.util.AnnotationUtil; import io.ebeaninternal.api.ConcurrencyMode; import io.ebeaninternal.api.SpiEbeanServer; import io.ebeaninternal.api.TransactionEventTable; +import io.ebeaninternal.server.cache.CacheChangeSet; import io.ebeaninternal.server.cache.SpiCacheManager; import io.ebeaninternal.server.core.InternString; import io.ebeaninternal.server.core.InternalConfiguration; @@ -449,21 +450,21 @@ public class BeanDescriptorManager implements BeanDescriptorMap { /** * For SQL based modifications we need to invalidate appropriate parts of the cache. */ - public void cacheNotify(TransactionEventTable.TableIUD tableIUD) { + public void cacheNotify(TransactionEventTable.TableIUD tableIUD, CacheChangeSet changeSet) { String tableName = tableIUD.getTableName().toLowerCase(); List> normalBeanTypes = tableToDescMap.get(tableName); if (normalBeanTypes != null) { // 'normal' entity beans based on a "base table" for (BeanDescriptor normalBeanType : normalBeanTypes) { - normalBeanType.cacheHandleBulkUpdate(tableIUD); + normalBeanType.cachePersistTableIUD(tableIUD, changeSet); } } List> viewBeans = tableToViewDescMap.get(tableName); if (viewBeans != null) { // entity beans based on a "view" for (BeanDescriptor viewBean : viewBeans) { - viewBean.cacheHandleBulkUpdate(tableIUD); + viewBean.cachePersistTableIUD(tableIUD, changeSet); } } } diff --git a/src/main/java/io/ebeaninternal/server/deploy/BeanPropertyAssocOne.java b/src/main/java/io/ebeaninternal/server/deploy/BeanPropertyAssocOne.java index 502a4bab1..8f8e3b66e 100644 --- a/src/main/java/io/ebeaninternal/server/deploy/BeanPropertyAssocOne.java +++ b/src/main/java/io/ebeaninternal/server/deploy/BeanPropertyAssocOne.java @@ -56,6 +56,7 @@ public class BeanPropertyAssocOne extends BeanPropertyAssoc implements STr private String deleteByParentIdInSql; private BeanPropertyAssocMany relationshipProperty; + private boolean cacheNotifyRelationship; /** * Create based on deploy information of an EmbeddedId. @@ -127,6 +128,13 @@ public class BeanPropertyAssocOne extends BeanPropertyAssoc implements STr } } + /** + * Derive late in lifecycle cache notification on this relationship. + */ + public void initialisePostTarget() { + this.cacheNotifyRelationship = isCacheNotifyRelationship(); + } + /** * Return the property value as an entity bean. */ @@ -141,25 +149,31 @@ public class BeanPropertyAssocOne extends BeanPropertyAssoc implements STr /** * Return true if this relationship needs to maintain/update L2 cache. */ - boolean isCacheNotify() { - return targetDescriptor.isBeanCaching() && relationshipProperty != null; + boolean isCacheNotifyRelationship() { + return relationshipProperty != null && targetDescriptor.isBeanCaching(); } /** * Clear the L2 relationship cache for this property. */ void cacheClear() { - if (isCacheNotify()) { + if (cacheNotifyRelationship) { targetDescriptor.cacheManyPropClear(relationshipProperty.getName()); } } + void cacheClear(CacheChangeSet changeSet) { + if (cacheNotifyRelationship) { + changeSet.addManyClear(targetDescriptor, relationshipProperty.getName()); + } + } + /** * Clear part of the L2 relationship cache for this property. */ void cacheDelete(boolean clear, EntityBean bean, CacheChangeSet changeSet) { - if (isCacheNotify()) { + if (cacheNotifyRelationship) { if (clear) { changeSet.addManyClear(targetDescriptor, relationshipProperty.getName()); } else { diff --git a/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java b/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java index a1b63873a..bb1690bfe 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java +++ b/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java @@ -4,6 +4,7 @@ import io.ebeaninternal.api.BinaryReadContext; import io.ebeaninternal.api.BinaryWritable; import io.ebeaninternal.api.BinaryWriteContext; import io.ebeaninternal.api.SpiEbeanServer; +import io.ebeaninternal.server.cache.CacheChangeSet; import io.ebeaninternal.server.core.PersistRequest; import io.ebeaninternal.server.deploy.BeanDescriptor; import io.ebeaninternal.server.deploy.id.IdBinder; @@ -126,11 +127,10 @@ public class BeanPersistIds implements BinaryWritable { /** * Notify the cache of this event that came from another server in the cluster. */ - void notifyCacheAndListener() { - // any change invalidates the query cache - beanDescriptor.clearQueryCache(); + public void notifyCache(CacheChangeSet changeSet) { + changeSet.addClearQuery(beanDescriptor); if (ids != null) { - beanDescriptor.cacheApplyInvalidate(ids); + changeSet.addBeanRemoveMany(beanDescriptor, ids); } } } diff --git a/src/main/java/io/ebeaninternal/server/transaction/DeleteByIdMap.java b/src/main/java/io/ebeaninternal/server/transaction/DeleteByIdMap.java index 5f76c4b03..ca2cc1f27 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/DeleteByIdMap.java +++ b/src/main/java/io/ebeaninternal/server/transaction/DeleteByIdMap.java @@ -29,7 +29,7 @@ public final class DeleteByIdMap { BeanDescriptor d = deleteIds.getBeanDescriptor(); List idValues = deleteIds.getIds(); if (idValues != null) { - d.cacheHandleDeleteByIds(idValues, changeSet); + d.cachePersistDeleteByIds(idValues, changeSet); } } } diff --git a/src/main/java/io/ebeaninternal/server/transaction/JdbcTransaction.java b/src/main/java/io/ebeaninternal/server/transaction/JdbcTransaction.java index 7fe53c827..bafd7a57f 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/JdbcTransaction.java +++ b/src/main/java/io/ebeaninternal/server/transaction/JdbcTransaction.java @@ -855,7 +855,7 @@ public class JdbcTransaction implements SpiTransaction, TxnProfileEventCodes { public TransactionEvent getEvent() { queryOnly = false; if (event == null) { - event = new TransactionEvent(startMillis); + event = new TransactionEvent(); } return event; } @@ -1052,7 +1052,7 @@ public class JdbcTransaction implements SpiTransaction, TxnProfileEventCodes { // the event has been sent to the transaction manager // for postCommit processing (l2 cache updates etc) // start a new transaction event - event = new TransactionEvent(startMillis); + event = new TransactionEvent(); } catch (Exception e) { doRollback(e); diff --git a/src/main/java/io/ebeaninternal/server/transaction/PostCommitProcessing.java b/src/main/java/io/ebeaninternal/server/transaction/PostCommitProcessing.java index 9797693fb..550bb8be4 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/PostCommitProcessing.java +++ b/src/main/java/io/ebeaninternal/server/transaction/PostCommitProcessing.java @@ -8,7 +8,6 @@ import io.ebeaninternal.api.TransactionEventTable.TableIUD; import io.ebeaninternal.server.cache.CacheChangeSet; import io.ebeaninternal.server.cluster.ClusterManager; import io.ebeaninternal.server.core.PersistRequestBean; -import io.ebeaninternal.server.deploy.BeanDescriptorManager; import io.ebeanservice.docstore.api.DocStoreUpdates; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -46,8 +45,6 @@ final class PostCommitProcessing { private final int txnDocStoreBatchSize; - private CacheChangeSet cacheChanges; - /** * Create for an external modification. */ @@ -83,31 +80,12 @@ final class PostCommitProcessing { } /** - * Notify the local part of L2 cache. + * Perform foreground cache notification if desired. */ void notifyLocalCache() { - processTableEvents(event.getEventTables()); if (manager.notifyL2CacheInForeground) { // process l2 cache changes in foreground - processCacheChanges(event.buildCacheChanges()); - } else { - // collect l2 cache changes for delayed background processing - cacheChanges = event.buildCacheChanges(); - } - } - - /** - * Table events are where SQL or external tools are used. In this case the - * cache is notified based on the table name (rather than bean type). - */ - private void processTableEvents(TransactionEventTable tableEvents) { - - if (tableEvents != null && !tableEvents.isEmpty()) { - // notify cache with table based changes - BeanDescriptorManager dm = manager.getBeanDescriptorManager(); - for (TableIUD tableIUD : tableEvents.values()) { - dm.cacheNotify(tableIUD); - } + processCacheChanges(); } } @@ -154,7 +132,9 @@ final class PostCommitProcessing { */ Runnable backgroundNotify() { return () -> { - processCacheChanges(cacheChanges); + if (!manager.notifyL2CacheInForeground) { + processCacheChanges(); + } localPersistListenersNotify(); notifyCluster(); processDocStoreUpdates(); @@ -164,7 +144,8 @@ final class PostCommitProcessing { /** * Apply the changes to the L2 caches. */ - private void processCacheChanges(CacheChangeSet cacheChanges) { + private void processCacheChanges() { + CacheChangeSet cacheChanges = event.buildCacheChanges(manager); if (cacheChanges != null) { Set touched = cacheChanges.touchedTables(); if (touched != null && !touched.isEmpty()) { diff --git a/src/main/java/io/ebeaninternal/server/transaction/TableModState.java b/src/main/java/io/ebeaninternal/server/transaction/TableModState.java index 443d35fd1..29b116aee 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/TableModState.java +++ b/src/main/java/io/ebeaninternal/server/transaction/TableModState.java @@ -38,7 +38,7 @@ public class TableModState implements QueryCacheEntryValidate, ServerCacheNotify tableModStamp.put(tableName, modTimestamp); } if (log.isDebugEnabled()) { - log.debug("TableModState updated - " + tableModStamp); + log.debug("TableModState updated - touched:{} modTimestamp:{}", touchedTables, modTimestamp); } } @@ -48,7 +48,10 @@ public class TableModState implements QueryCacheEntryValidate, ServerCacheNotify boolean isValid(Set tables, long sinceTimestamp) { for (String tableName : tables) { Long modTime = tableModStamp.get(tableName); - if (modTime != null && modTime > sinceTimestamp ) { + if (modTime != null && modTime >= sinceTimestamp ) { + if (log.isTraceEnabled()) { + log.trace("Invalidate on table:{}", tableName); + } return false; } } @@ -75,6 +78,9 @@ public class TableModState implements QueryCacheEntryValidate, ServerCacheNotify // use local clock - for slightly more aggressive invalidation (as later) // that removes any concern regarding clock syncing across cluster + if (log.isDebugEnabled()) { + log.debug("ServerCacheNotification:{}", notification); + } touch(notification.getDependentTables(), clockService.nowMillis()); } @@ -88,6 +94,9 @@ public class TableModState implements QueryCacheEntryValidate, ServerCacheNotify // use local clock - for slightly more aggressive invalidation (as later) // that removes any concern regarding clock syncing across cluster + if (log.isDebugEnabled()) { + log.debug("RemoteTableMod:{}", tableMod); + } touch(tableMod.getTables(), clockService.nowMillis()); } } diff --git a/src/main/java/io/ebeaninternal/server/transaction/TransactionManager.java b/src/main/java/io/ebeaninternal/server/transaction/TransactionManager.java index f330e8cce..147839461 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/TransactionManager.java +++ b/src/main/java/io/ebeaninternal/server/transaction/TransactionManager.java @@ -28,6 +28,7 @@ import io.ebeaninternal.api.TransactionEventTable.TableIUD; import io.ebeaninternal.metric.MetricFactory; import io.ebeaninternal.metric.TimedMetric; import io.ebeaninternal.metric.TimedMetricMap; +import io.ebeaninternal.server.cache.CacheChangeSet; import io.ebeaninternal.server.cluster.ClusterManager; import io.ebeaninternal.server.core.ClockService; import io.ebeaninternal.server.deploy.BeanDescriptorManager; @@ -451,7 +452,7 @@ public class TransactionManager implements SpiTransactionManager { private void externalModificationEvent(TransactionEventTable tableEvents) { - TransactionEvent event = new TransactionEvent(clockNowMillis()); + TransactionEvent event = new TransactionEvent(); event.add(tableEvents); PostCommitProcessing postCommit = new PostCommitProcessing(clusterManager, this, event); @@ -468,15 +469,17 @@ public class TransactionManager implements SpiTransactionManager { clusterLogger.debug("processing {}", remoteEvent); } + CacheChangeSet changeSet = new CacheChangeSet(clockNowMillis()); + RemoteTableMod tableMod = remoteEvent.getRemoteTableMod(); if (tableMod != null) { - tableModState.notify(tableMod); + changeSet.addInvalidate(tableMod.getTables()); } List tableIUDList = remoteEvent.getTableIUDList(); if (tableIUDList != null) { for (TableIUD tableIUD : tableIUDList) { - beanDescriptorManager.cacheNotify(tableIUD); + beanDescriptorManager.cacheNotify(tableIUD, changeSet); } } @@ -484,10 +487,12 @@ public class TransactionManager implements SpiTransactionManager { // processes both Bean IUD and DeleteById List beanPersistList = remoteEvent.getBeanPersistList(); if (beanPersistList != null) { - for (BeanPersistIds aBeanPersistList : beanPersistList) { - aBeanPersistList.notifyCacheAndListener(); + for (BeanPersistIds persistIds : beanPersistList) { + persistIds.notifyCache(changeSet); } } + + changeSet.apply(); } /** diff --git a/src/test/java/org/tests/cache/TestCacheCollectionIds.java b/src/test/java/org/tests/cache/TestCacheCollectionIds.java index 515fb901f..bb09dd70e 100644 --- a/src/test/java/org/tests/cache/TestCacheCollectionIds.java +++ b/src/test/java/org/tests/cache/TestCacheCollectionIds.java @@ -23,9 +23,13 @@ import org.tests.model.basic.ResetBasicData; import java.util.ArrayList; import java.util.List; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + public class TestCacheCollectionIds extends BaseTestCase { - ServerCacheManager cacheManager = Ebean.getServerCacheManager(); + private ServerCacheManager cacheManager = Ebean.getServerCacheManager(); @Test public void test() { @@ -45,7 +49,7 @@ public class TestCacheCollectionIds extends BaseTestCase { List list = Ebean.find(Customer.class).setAutoTune(false).setBeanCacheMode(CacheMode.PUT) .order().asc("id").findList(); - Assert.assertTrue(list.size() > 1); + assertTrue(list.size() > 1); // Assert.assertEquals(list.size(), // custCache.getStatistics(false).getSize()); @@ -53,7 +57,7 @@ public class TestCacheCollectionIds extends BaseTestCase { List contacts = customer.getContacts(); // Assert.assertEquals(0, custManyIdsCache.getStatistics(false).getSize()); contacts.size(); - Assert.assertTrue(contacts.size() > 1); + assertTrue(contacts.size() > 1); // Assert.assertEquals(1, custManyIdsCache.getStatistics(false).getSize()); // Assert.assertEquals(0, // custManyIdsCache.getStatistics(false).getHitCount()); @@ -77,7 +81,7 @@ public class TestCacheCollectionIds extends BaseTestCase { awaitL2Cache(); int currentNumContacts2 = fetchCustomer(customer.getId()); - Assert.assertEquals(currentNumContacts + 1, currentNumContacts2); + assertEquals(currentNumContacts + 1, currentNumContacts2); // cleanup Ebean.delete(newContact); @@ -119,8 +123,8 @@ public class TestCacheCollectionIds extends BaseTestCase { CachedManyIds cachedManyIds = (CachedManyIds) cachedBeanCountriesCache.get(cachedBean.getId()); // confirm the starting data and cache entry - Assert.assertEquals(2, dummyToLoad.getCountries().size()); - Assert.assertEquals(2, cachedManyIds.getIdList().size()); + assertEquals(2, dummyToLoad.getCountries().size()); + assertEquals(2, cachedManyIds.getIdList().size()); // act @@ -136,10 +140,10 @@ public class TestCacheCollectionIds extends BaseTestCase { cachedManyIds = (CachedManyIds) cachedBeanCountriesCache.get(result.getId()); // assert that data and cache both show correct data - Assert.assertEquals(1, result.getCountries().size()); - Assert.assertEquals(1, cachedManyIds.getIdList().size()); - Assert.assertFalse(cachedManyIds.getIdList().contains("NZ")); - Assert.assertTrue(cachedManyIds.getIdList().contains("AU")); + assertEquals(1, result.getCountries().size()); + assertEquals(1, cachedManyIds.getIdList().size()); + assertFalse(cachedManyIds.getIdList().contains("NZ")); + assertTrue(cachedManyIds.getIdList().contains("AU")); } @@ -167,8 +171,8 @@ public class TestCacheCollectionIds extends BaseTestCase { CachedManyIds cachedManyIds = (CachedManyIds) cachedBeanCountriesCache.get(cachedBean.getId()); // confirm the starting data and cache entry - Assert.assertEquals(2, dummyToLoad.getCountries().size()); - Assert.assertEquals(2, cachedManyIds.getIdList().size()); + assertEquals(2, dummyToLoad.getCountries().size()); + assertEquals(2, cachedManyIds.getIdList().size()); // act - this time update the name property so the bean is dirty @@ -185,10 +189,10 @@ public class TestCacheCollectionIds extends BaseTestCase { cachedManyIds = (CachedManyIds) cachedBeanCountriesCache.get(result.getId()); // assert that data and cache both show correct data - Assert.assertEquals(1, result.getCountries().size()); - Assert.assertEquals(1, cachedManyIds.getIdList().size()); - Assert.assertFalse(cachedManyIds.getIdList().contains("NZ")); - Assert.assertTrue(cachedManyIds.getIdList().contains("AU")); + assertEquals(1, result.getCountries().size()); + assertEquals(1, cachedManyIds.getIdList().size()); + assertFalse(cachedManyIds.getIdList().contains("NZ")); + assertTrue(cachedManyIds.getIdList().contains("AU")); } /** @@ -209,20 +213,20 @@ public class TestCacheCollectionIds extends BaseTestCase { // clear the cache ServerCache cachedBeanCountriesCache = cacheManager.getCollectionIdsCache(OCachedBean.class, "countries"); cachedBeanCountriesCache.clear(); - Assert.assertEquals(0, cachedBeanCountriesCache.size()); + assertEquals(0, cachedBeanCountriesCache.size()); // load the cache OCachedBean dummyLoad = Ebean.find(OCachedBean.class, cachedBean.getId()); List dummyCountries = dummyLoad.getCountries(); - Assert.assertEquals(2, dummyCountries.size()); + assertEquals(2, dummyCountries.size()); // assert that the cache contains the expected entry - Assert.assertEquals("countries cache now loaded with 1 entry", 1, cachedBeanCountriesCache.size()); + assertEquals("countries cache now loaded with 1 entry", 1, cachedBeanCountriesCache.size()); CachedManyIds dummyEntry = (CachedManyIds) cachedBeanCountriesCache.get(dummyLoad.getId()); Assert.assertNotNull(dummyEntry); - Assert.assertEquals("2 ids in the entry", 2, dummyEntry.getIdList().size()); - Assert.assertTrue(dummyEntry.getIdList().contains("NZ")); - Assert.assertTrue(dummyEntry.getIdList().contains("AU")); + assertEquals("2 ids in the entry", 2, dummyEntry.getIdList().size()); + assertTrue(dummyEntry.getIdList().contains("NZ")); + assertTrue(dummyEntry.getIdList().contains("AU")); // act - this should invalidate our cache entry @@ -234,18 +238,18 @@ public class TestCacheCollectionIds extends BaseTestCase { Ebean.update(update); awaitL2Cache(); - Assert.assertEquals("countries entry still there (but updated)", 1, cachedBeanCountriesCache.size()); + assertEquals("countries entry still there (but updated)", 1, cachedBeanCountriesCache.size()); CachedManyIds cachedManyIds = (CachedManyIds) cachedBeanCountriesCache.get(update.getId()); // assert cache updated - Assert.assertEquals(1, cachedManyIds.getIdList().size()); - Assert.assertFalse(cachedManyIds.getIdList().contains("NZ")); - Assert.assertTrue(cachedManyIds.getIdList().contains("AU")); + assertEquals(1, cachedManyIds.getIdList().size()); + assertFalse(cachedManyIds.getIdList().contains("NZ")); + assertTrue(cachedManyIds.getIdList().contains("AU")); // assert countries good OCachedBean result = Ebean.find(OCachedBean.class, cachedBean.getId()); - Assert.assertEquals(1, result.getCountries().size()); + assertEquals(1, result.getCountries().size()); } @@ -300,7 +304,7 @@ public class TestCacheCollectionIds extends BaseTestCase { Update update = Ebean.createUpdate(OrderDetail.class, updStatement); update.set("id", orderDetail1.getId()); int rows = update.execute(); - Assert.assertEquals(1, rows); + assertEquals(1, rows); // read the order from cache Order orderFromCache = Ebean.find(Order.class, 1L); @@ -335,7 +339,7 @@ public class TestCacheCollectionIds extends BaseTestCase { SqlUpdate update = Ebean.createSqlUpdate(updStatement); update.setParameter("id", orderDetail1.getId()); int rows = update.execute(); - Assert.assertEquals(1, rows); + assertEquals(1, rows); // We need to notify the cache manually Ebean.externalModification("o_order_detail", false, false, true); diff --git a/src/test/java/org/tests/cache/TestCacheDelete.java b/src/test/java/org/tests/cache/TestCacheDelete.java index 9af7aeae4..a5f5e111e 100644 --- a/src/test/java/org/tests/cache/TestCacheDelete.java +++ b/src/test/java/org/tests/cache/TestCacheDelete.java @@ -2,10 +2,11 @@ package org.tests.cache; import io.ebean.BaseTestCase; import io.ebean.Ebean; +import org.junit.Test; import org.tests.model.basic.OCachedBean; import org.tests.model.basic.OCachedBeanChild; -import org.junit.Assert; -import org.junit.Test; + +import static org.junit.Assert.assertEquals; /** * Test class testing deleting/invalidating of cached beans @@ -27,7 +28,7 @@ public class TestCacheDelete extends BaseTestCase { Ebean.save(parentBean); // confirm there are 2 children loaded from the parent - Assert.assertEquals(2, Ebean.find(OCachedBean.class, parentBean.getId()).getChildren().size()); + assertEquals(2, Ebean.find(OCachedBean.class, parentBean.getId()).getChildren().size()); // ensure cache has been populated Ebean.find(OCachedBeanChild.class, child.getId()); @@ -39,6 +40,6 @@ public class TestCacheDelete extends BaseTestCase { awaitL2Cache(); OCachedBean beanFromCache = Ebean.find(OCachedBean.class, parentBean.getId()); - Assert.assertEquals(1, beanFromCache.getChildren().size()); + assertEquals(1, beanFromCache.getChildren().size()); } } diff --git a/src/test/java/org/tests/cache/TestQueryCacheTableDependency.java b/src/test/java/org/tests/cache/TestQueryCacheTableDependency.java index f578fdfdf..494d5c377 100644 --- a/src/test/java/org/tests/cache/TestQueryCacheTableDependency.java +++ b/src/test/java/org/tests/cache/TestQueryCacheTableDependency.java @@ -5,6 +5,7 @@ import io.ebean.Ebean; import io.ebean.cache.ServerCache; import org.junit.Test; import org.tests.model.basic.Address; +import org.tests.model.basic.Contact; import org.tests.model.basic.Customer; import org.tests.model.basic.ResetBasicData; @@ -46,7 +47,7 @@ public class TestQueryCacheTableDependency extends BaseTestCase { .where().eq("billingAddress.line2", "St Lukes") .findCount(); - assertThat(custs).isEqualTo(2); // cache says 3 + assertThat(custs).isEqualTo(2); custs = Ebean.find(Customer.class).setUseQueryCache(true).setReadOnly(true) @@ -55,5 +56,70 @@ public class TestQueryCacheTableDependency extends BaseTestCase { assertThat(custs).isEqualTo(1); + Ebean.update(Address.class) + .set("line2", "St Lucky2") + .where().eq("line2", "St Lucky") + .update(); + + custs = Ebean.find(Customer.class).setUseQueryCache(true).setReadOnly(true) + .where().eq("billingAddress.line2", "St Lucky") + .findCount(); + + assertThat(custs).isEqualTo(0); + + custs = Ebean.find(Customer.class).setUseQueryCache(true).setReadOnly(true) + .where().eq("billingAddress.line2", "St Lucky2") + .findCount(); + assertThat(custs).isEqualTo(1); + + Ebean.createSqlUpdate("update O_ADDRESS set line_2=? where line_2=?") + .setNextParameter("St Lucky3") + .setNextParameter("St Lucky2") + .execute(); + + custs = Ebean.find(Customer.class).setUseQueryCache(true).setReadOnly(true) + .where().eq("billingAddress.line2", "St Lucky2") + .findCount(); + assertThat(custs).isEqualTo(0); + + custs = Ebean.find(Customer.class).setUseQueryCache(true).setReadOnly(true) + .where().eq("billingAddress.line2", "St Lucky3") + .findCount(); + + assertThat(custs).isEqualTo(1); + + } + + @Test + public void testFindCountOnOtherL2Cached() { + + ResetBasicData.reset(); + + Customer fi = Ebean.find(Customer.class).where().eq("name", "Fiona").findOne(); + + int custCount0 = Ebean.find(Customer.class).setUseQueryCache(true).setReadOnly(true) + .where() + .eq("name", "Fiona") + .isNull("contacts.phone") + .findCount(); + + assertThat(custCount0).isEqualTo(1); + + int updateRows = Ebean.update(Contact.class) + .set("phone", "1234") + .where() + .eq("customer.id", fi.getId()) + .update(); + + assertThat(updateRows).isGreaterThan(0); + + int custCount1 = Ebean.find(Customer.class).setUseQueryCache(true).setReadOnly(true) + .where() + .eq("name", "Fiona") + .isNull("contacts.phone") + .findCount(); + + assertThat(custCount1).isEqualTo(0); + } } diff --git a/src/test/resources/logback-test.xml b/src/test/resources/logback-test.xml index 7ab47450e..214c25900 100644 --- a/src/test/resources/logback-test.xml +++ b/src/test/resources/logback-test.xml @@ -86,6 +86,8 @@ + +