From 3f5009aa91b5ad9b009edb8fa392fc853a2a4409 Mon Sep 17 00:00:00 2001 From: Roland Praml Date: Mon, 11 Oct 2021 16:25:21 +0200 Subject: [PATCH] Use WeakReferences when in iterate query mode --- .../io/ebean/bean/PersistenceContext.java | 9 +- .../server/core/OrmQueryRequest.java | 15 +- .../DefaultPersistenceContext.java | 131 +++++++-- .../transaction/WeakPersistenceContext.java | 272 ------------------ .../tests/basic/TestPersistenceContext.java | 66 +++++ 5 files changed, 195 insertions(+), 298 deletions(-) delete mode 100644 ebean-core/src/main/java/io/ebeaninternal/server/transaction/WeakPersistenceContext.java diff --git a/ebean-api/src/main/java/io/ebean/bean/PersistenceContext.java b/ebean-api/src/main/java/io/ebean/bean/PersistenceContext.java index 33215595a..ed03080cd 100644 --- a/ebean-api/src/main/java/io/ebean/bean/PersistenceContext.java +++ b/ebean-api/src/main/java/io/ebean/bean/PersistenceContext.java @@ -60,10 +60,15 @@ public interface PersistenceContext { int size(Class rootType); /** - * Return a copy of the Persistence context to use for large query iteration. + * Signalizes the PersistenceContext, the begin for large query iteration. */ - PersistenceContext forIterate(); + void beginIterate(); + /** + * Signalizes the PersistenceContext, the end for large query iteration. + */ + void endIterate(); + /** * Wrapper on a bean to also indicate if a bean has been deleted. *

diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java b/ebean-core/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java index a56319fd4..6ef1e352b 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java @@ -225,6 +225,9 @@ public final class OrmQueryRequest extends BeanRequest implements SpiOrmQuery createdTransaction = true; } persistenceContext = persistenceContext(query, transaction); + if (Type.ITERATE == query.getType()) { + persistenceContext.beginIterate(); + } loadContext = new DLoadContext(this, secondaryQueries); } @@ -233,6 +236,9 @@ public final class OrmQueryRequest extends BeanRequest implements SpiOrmQuery */ @Override public void rollbackTransIfRequired() { + if (Type.ITERATE == query.getType()) { + persistenceContext.endIterate(); + } if (createdTransaction) { try { transaction.end(); @@ -278,11 +284,7 @@ public final class OrmQueryRequest extends BeanRequest implements SpiOrmQuery if (scope == PersistenceContextScope.QUERY || t == null) { return new DefaultPersistenceContext(); } - if (Type.ITERATE == query.getType()) { - return t.getPersistenceContext().forIterate(); - } else { - return t.getPersistenceContext(); - } + return t.getPersistenceContext(); } /** @@ -292,6 +294,9 @@ public final class OrmQueryRequest extends BeanRequest implements SpiOrmQuery */ @Override public void endTransIfRequired() { + if (Type.ITERATE == query.getType()) { + persistenceContext.endIterate(); + } if (createdTransaction && transaction.isActive()) { transaction.commit(); if (query.getType().isUpdate()) { diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/transaction/DefaultPersistenceContext.java b/ebean-core/src/main/java/io/ebeaninternal/server/transaction/DefaultPersistenceContext.java index 911c54c47..1ff3d44fa 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/transaction/DefaultPersistenceContext.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/transaction/DefaultPersistenceContext.java @@ -6,6 +6,9 @@ import io.ebeaninternal.api.SpiBeanType; import io.ebeaninternal.api.SpiBeanTypeManager; import io.ebeaninternal.api.SpiPersistenceContext; +import java.lang.ref.Reference; +import java.lang.ref.ReferenceQueue; +import java.lang.ref.WeakReference; import java.util.*; import java.util.concurrent.locks.ReentrantLock; @@ -27,6 +30,14 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { private final HashMap, ClassContext> typeCache = new HashMap<>(); private final ReentrantLock lock = new ReentrantLock(); + private final ReferenceQueue queue = new ReferenceQueue(); + + /** + * When we are inside an iterate loop, we will add only WeakReferences. This + * allows the JVM GC to collect beans, which are not referenced elsewhere. In + * normal operation, we will use hard references, to avoid performance impact + */ + private int iterateDepth; /** * Create a new PersistenceContext. @@ -34,16 +45,25 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { public DefaultPersistenceContext() { } - /** - * Return a PersistenceContext to use for streaming queries. - */ @Override - public PersistenceContext forIterate() { - WeakPersistenceContext weak = new WeakPersistenceContext(); - for (ClassContext value : typeCache.values()) { - weak.add(value.rootType, value.deleteSet, value.map); + public void beginIterate() { + lock.lock(); + try { + iterateDepth++; + } finally { + lock.unlock(); + } + } + + @Override + public void endIterate() { + lock.lock(); + try { + iterateDepth--; + expungeStaleEntries(); // when leaving the iterator, cleanup. + } finally { + lock.unlock(); } - return weak; } /** @@ -53,7 +73,8 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { public void put(Class rootType, Object id, Object bean) { lock.lock(); try { - classContext(rootType).put(id, bean); + expungeStaleEntries(); + classContext(rootType).useReferences(iterateDepth > 0).put(id, bean); } finally { lock.unlock(); } @@ -63,7 +84,8 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { public Object putIfAbsent(Class rootType, Object id, Object bean) { lock.lock(); try { - return classContext(rootType).putIfAbsent(id, bean); + expungeStaleEntries(); + return classContext(rootType).useReferences(iterateDepth > 0).putIfAbsent(id, bean); } finally { lock.unlock(); } @@ -76,6 +98,7 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { public Object get(Class rootType, Object id) { lock.lock(); try { + expungeStaleEntries(); return classContext(rootType).get(id); } finally { lock.unlock(); @@ -86,6 +109,7 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { public WithOption getWithOption(Class rootType, Object id) { lock.lock(); try { + expungeStaleEntries(); return classContext(rootType).getWithOption(id); } finally { lock.unlock(); @@ -99,6 +123,7 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { public int size(Class rootType) { lock.lock(); try { + expungeStaleEntries(); ClassContext classMap = typeCache.get(rootType); return classMap == null ? 0 : classMap.size(); } finally { @@ -114,6 +139,7 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { lock.lock(); try { typeCache.clear(); + expungeStaleEntries(); } finally { lock.unlock(); } @@ -127,6 +153,7 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { if (classMap != null) { classMap.clear(); } + expungeStaleEntries(); } finally { lock.unlock(); } @@ -140,6 +167,7 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { if (classMap != null && id != null) { classMap.deleted(id); } + expungeStaleEntries(); } finally { lock.unlock(); } @@ -153,6 +181,7 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { if (classMap != null && id != null) { classMap.remove(id); } + expungeStaleEntries(); } finally { lock.unlock(); } @@ -162,6 +191,7 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { public List dirtyBeans(SpiBeanTypeManager manager) { lock.lock(); try { + expungeStaleEntries(); List list = new ArrayList<>(); for (ClassContext classContext : typeCache.values()) { classContext.dirtyBeans(manager, list); @@ -172,10 +202,23 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { } } + /** + * When there is a queue, poll it to remove stale entries from the map. Note: + * This is always done AFTER useReferences was called with + * true. Polling an empty queue has no performance impact. + */ + private void expungeStaleEntries() { + Reference ref; + while ((ref = queue.poll()) != null) { + ((BeanRef)ref).expunge(); + } + } + @Override public String toString() { lock.lock(); try { + expungeStaleEntries(); return typeCache.toString(); } finally { lock.unlock(); @@ -183,26 +226,47 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { } private ClassContext classContext(Class rootType) { - return typeCache.computeIfAbsent(rootType, k -> new ClassContext(rootType)); + return typeCache.computeIfAbsent(rootType, k -> new ClassContext(k, queue)); } - + private static class ClassContext { private final Map map = new HashMap<>(); private final Class rootType; + private final ReferenceQueue queue; private Set deleteSet; + + private boolean useReferences; + private int weakCount; - private ClassContext(Class rootType) { + + private ClassContext(Class rootType, ReferenceQueue queue) { this.rootType = rootType; + this.queue = queue; } + /** + * When called with "true", initialize referenceQueue and store BeanRefs instead + * of real object references. + */ + public ClassContext useReferences(boolean useReferences) { + this.useReferences = useReferences; + return this; + } + + @Override public String toString() { - return "size:" + map.size(); + return "size:" + map.size() + " (" + weakCount + " weak)"; } private Object get(Object id) { - return map.get(id); + Object ret = map.get(id); + if (ret instanceof BeanRef) { + return ((BeanRef)ret).get(); + } else { + return ret; + } } private WithOption getWithOption(Object id) { @@ -220,12 +284,17 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { return existingValue; } // put the new value and return null indicating the put was successful - map.put(id, bean); + put(id, bean); return null; } private void put(Object id, Object b) { - map.put(id, b); + if (useReferences) { + weakCount++; + map.put(id, new BeanRef(this, id, b, queue)); + } else { + map.put(id, b); + } } private int size() { @@ -234,10 +303,14 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { private void clear() { map.clear(); + weakCount = 0; } private void remove(Object id) { - map.remove(id); + Object ret = map.remove(id); + if (ret instanceof BeanRef) { + weakCount--; + } } private void deleted(Object id) { @@ -245,7 +318,7 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { deleteSet = new HashSet<>(); } deleteSet.add(id); - map.remove(id); + remove(id); } /** @@ -254,6 +327,10 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { void dirtyBeans(SpiBeanTypeManager manager, List list) { final SpiBeanType beanType = manager.beanType(rootType); for (Object value : map.values()) { + if (queue != null && value instanceof BeanRef) { + value = ((BeanRef) value).get(); + if (value == null) continue; + } EntityBean bean = (EntityBean) value; if (bean._ebean_getIntercept().isDirty() || beanType.isToManyDirty(bean)) { list.add(value); @@ -262,4 +339,20 @@ public final class DefaultPersistenceContext implements SpiPersistenceContext { } } + private static class BeanRef extends WeakReference { + + private final ClassContext classContext; + private final Object key; + + private BeanRef(ClassContext classContext, Object key, Object referent, ReferenceQueue q) { + super(referent, q); + this.classContext = classContext; + this.key = key; + } + + private void expunge() { + classContext.remove(key); + } + + } } diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/transaction/WeakPersistenceContext.java b/ebean-core/src/main/java/io/ebeaninternal/server/transaction/WeakPersistenceContext.java deleted file mode 100644 index 4d7371bf9..000000000 --- a/ebean-core/src/main/java/io/ebeaninternal/server/transaction/WeakPersistenceContext.java +++ /dev/null @@ -1,272 +0,0 @@ -package io.ebeaninternal.server.transaction; - -import io.ebean.bean.EntityBean; -import io.ebean.bean.PersistenceContext; -import io.ebeaninternal.api.SpiBeanType; -import io.ebeaninternal.api.SpiBeanTypeManager; -import io.ebeaninternal.api.SpiPersistenceContext; - -import java.lang.ref.Reference; -import java.lang.ref.ReferenceQueue; -import java.lang.ref.WeakReference; -import java.util.*; -import java.util.concurrent.locks.ReentrantLock; - -/** - * Weak reference based PersistenceContext for "streaming queries". - */ -final class WeakPersistenceContext implements SpiPersistenceContext { - - private final HashMap, WeakClassContext> typeCache = new HashMap<>(); - private final ReentrantLock lock = new ReentrantLock(); - - WeakPersistenceContext() { - } - - /** - * Load from the initiating persistence context. - */ - void add(Class rootType, Set deleteSet, Map map) { - typeCache.put(rootType, new WeakClassContext( rootType, deleteSet, map)); - } - - @Override - public PersistenceContext forIterate() { - return this; - } - - @Override - public void put(Class rootType, Object id, Object bean) { - lock.lock(); - try { - classContext(rootType).put(id, bean); - } finally { - lock.unlock(); - } - } - - @Override - public Object putIfAbsent(Class rootType, Object id, Object bean) { - lock.lock(); - try { - return classContext(rootType).putIfAbsent(id, bean); - } finally { - lock.unlock(); - } - } - - @Override - public Object get(Class rootType, Object id) { - lock.lock(); - try { - return classContext(rootType).get(id); - } finally { - lock.unlock(); - } - } - - @Override - public WithOption getWithOption(Class rootType, Object id) { - lock.lock(); - try { - return classContext(rootType).getWithOption(id); - } finally { - lock.unlock(); - } - } - - @Override - public int size(Class rootType) { - lock.lock(); - try { - WeakClassContext classMap = typeCache.get(rootType); - return classMap == null ? 0 : classMap.size(); - } finally { - lock.unlock(); - } - } - - @Override - public void clear() { - lock.lock(); - try { - typeCache.clear(); - } finally { - lock.unlock(); - } - } - - @Override - public void clear(Class rootType) { - lock.lock(); - try { - WeakClassContext classMap = typeCache.get(rootType); - if (classMap != null) { - classMap.clear(); - } - } finally { - lock.unlock(); - } - } - - @Override - public void deleted(Class rootType, Object id) { - lock.lock(); - try { - WeakClassContext classMap = typeCache.get(rootType); - if (classMap != null && id != null) { - classMap.deleted(id); - } - } finally { - lock.unlock(); - } - } - - @Override - public void clear(Class rootType, Object id) { - lock.lock(); - try { - WeakClassContext classMap = typeCache.get(rootType); - if (classMap != null && id != null) { - classMap.remove(id); - } - } finally { - lock.unlock(); - } - } - - @Override - public List dirtyBeans(SpiBeanTypeManager manager) { - lock.lock(); - try { - List list = new ArrayList<>(); - for (WeakClassContext classContext : typeCache.values()) { - classContext.dirtyBeans(manager, list); - } - return list; - } finally { - lock.unlock(); - } - } - - @Override - public String toString() { - lock.lock(); - try { - return typeCache.toString(); - } finally { - lock.unlock(); - } - } - - private WeakClassContext classContext(Class rootType) { - return typeCache.computeIfAbsent(rootType, k -> new WeakClassContext(rootType)); - } - - private static class WeakClassContext { - - private final ReferenceQueue queue = new ReferenceQueue<>(); - private final Map map = new HashMap<>(); - private final Class rootType; - private Set deleteSet; - - private WeakClassContext(Class rootType) { - this.rootType = rootType; - } - - private WeakClassContext(Class rootType, Set initialDeleteSet, Map initialMap) { - this.rootType = rootType; - this.deleteSet = initialDeleteSet == null ? null : new HashSet<>(initialDeleteSet); - for (Map.Entry entry : initialMap.entrySet()) { - Object id = entry.getKey(); - map.put(id, new BeanRef(id, entry.getValue(), queue)); - } - } - - private void expungeStaleEntries() { - Reference ref; - while ((ref = queue.poll()) != null) { - map.remove(((BeanRef)ref).key()); - } - } - - @Override - public String toString() { - return "size:" + map.size(); - } - - private Object get(Object id) { - expungeStaleEntries(); - Reference ref = map.get(id); - return ref == null ? null : ref.get(); - } - - private WithOption getWithOption(Object id) { - if (deleteSet != null && deleteSet.contains(id)) { - return WithOption.DELETED; - } - Object bean = get(id); - return (bean == null) ? null : new WithOption(bean); - } - - private Object putIfAbsent(Object id, Object bean) { - Object existingValue = get(id); - if (existingValue != null) { - // it is not absent - return existingValue; - } - // put the new value and return null indicating the put was successful - put(id, bean); - return null; - } - - private void put(Object id, Object b) { - expungeStaleEntries(); - map.put(id, new BeanRef(id, b, queue)); - } - - private int size() { - return map.size(); - } - - private void clear() { - map.clear(); - } - - private void remove(Object id) { - map.remove(id); - } - - private void deleted(Object id) { - if (deleteSet == null) { - deleteSet = new HashSet<>(); - } - deleteSet.add(id); - map.remove(id); - } - - /** - * Add the dirty beans to the list. - */ - void dirtyBeans(SpiBeanTypeManager manager, List list) { - final SpiBeanType beanType = manager.beanType(rootType); - for (BeanRef value : map.values()) { - EntityBean bean = (EntityBean) value.get(); - if (bean != null && (bean._ebean_getIntercept().isDirty() || beanType.isToManyDirty(bean))) { - list.add(bean); - } - } - } - } - - private static class BeanRef extends WeakReference { - private final Object key; - private BeanRef(Object key, Object referent, ReferenceQueue q) { - super(referent, q); - this.key = key; - } - Object key() { - return key; - } - } -} diff --git a/ebean-test/src/test/java/org/tests/basic/TestPersistenceContext.java b/ebean-test/src/test/java/org/tests/basic/TestPersistenceContext.java index ea894091a..4d79e5b89 100644 --- a/ebean-test/src/test/java/org/tests/basic/TestPersistenceContext.java +++ b/ebean-test/src/test/java/org/tests/basic/TestPersistenceContext.java @@ -2,16 +2,23 @@ package org.tests.basic; import io.ebean.BaseTestCase; import io.ebean.DB; +import io.ebean.Query; +import io.ebean.Transaction; +import io.ebeaninternal.api.SpiPersistenceContext; +import io.ebeaninternal.api.SpiTransaction; + import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.tests.model.basic.Customer; import org.tests.model.basic.Order; import org.tests.model.basic.ResetBasicData; +import java.util.List; import java.util.WeakHashMap; import java.util.concurrent.atomic.AtomicInteger; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertTrue; public class TestPersistenceContext extends BaseTestCase { @@ -112,6 +119,65 @@ public class TestPersistenceContext extends BaseTestCase { } }); } + } + + @Disabled // run manually + @Test + void testPcScopes() throws InterruptedException { + for (int i = 0; i < 5000; i++) { + Customer c = new Customer(); + c.setName("Customer #" + i); + DB.save(c); + Order o = new Order(); + o.setCustomer(c); + DB.save(o); + } + + try (Transaction txn = DB.beginTransaction()) { + List first100 = DB.find(Customer.class).where().le("id", 100).findList(); + assertEquals(100, first100.size()); + for (Customer c : first100) { + c.setSmallnote("one of the first 100"); + } + + Customer[] lastBean = new Customer[1]; + DB.find(Customer.class).setLazyLoadBatchSize(1).findEach(customer -> { + if (customer.getId() <= 100) { + assertEquals("one of the first 100", customer.getSmallnote()); + // nested finds + DB.find(Order.class).where().eq("customer", customer).findEach(20, consumer->{}); + } else { + assertNotEquals("one of the first 100", customer.getSmallnote()); + } + lastBean[0] = customer; + }); + + SpiPersistenceContext pc = ((SpiTransaction) txn).getPersistenceContext(); + System.out.println(pc.toString()); // Customer=size:5000 (4900 weak), Order=size:100 (100 weak) + + System.gc(); + Thread.sleep(100); + pc.get(Customer.class, 1); // trigger expungeStaleEntries + pc.get(Order.class, 1); // Checkme: Would it be a problem, that every classContext has its own queue + System.out.println(pc.toString()); // Customer=size:101 (1 weak), no orders + + first100 = DB.find(Customer.class).where().le("id", 100).findList(); + for (Customer c : first100) { + assertEquals("one of the first 100", c.getSmallnote()); + } + Customer lastBeanFromDb = DB.find(Customer.class).setId(lastBean[0].getId()).findOne(); + assertTrue(lastBeanFromDb == lastBean[0]); + + // read 200 + DB.find(Customer.class).where().le("id", 200).findList(); + System.out.println(pc.toString()); // Customer=size:201 (1 weak) + + lastBean[0] = null; + lastBeanFromDb = null; + System.gc(); + Thread.sleep(100); + System.out.println(pc.toString()); // Customer=size:200 (0 weak) + } } }