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 super Object> 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 super Object> 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)
+ }
}
}