#1823 - ENH: Add query.findStream() and query.findLargeStream()

This commit is contained in:
rob bygrave
2019-10-10 13:44:29 +13:00
parent 5d433ff603
commit fdc67b3ce5
12 changed files with 249 additions and 13 deletions
@@ -13,6 +13,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.function.Consumer;
import java.util.function.Predicate;
import java.util.stream.Stream;
/**
* The extended API for Database.
@@ -107,6 +108,24 @@ public interface ExtendedServer {
@Nonnull
<T> QueryIterator<T> findIterate(Query<T> query, Transaction transaction);
/**
* Return the query result as a Stream using a single persistence context.
* <p>
* Note that the stream needs to be closed so use with try with resources.
* </p>
*/
@Nonnull
<T> Stream<T> findStream(Query<T> query, Transaction transaction);
/**
* Return the query result as a Stream (with multiple persistence contexts).
* <p>
* Note that the stream needs to be closed so use with try with resources.
* </p>
*/
@Nonnull
<T> Stream<T> findLargeStream(Query<T> query, Transaction transaction);
/**
* Execute the query visiting the each bean one at a time.
* <p>
+44
View File
@@ -11,6 +11,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.function.Consumer;
import java.util.function.Predicate;
import java.util.stream.Stream;
/**
* Object relational query for finding a List, Set, Map or single entity bean.
@@ -702,6 +703,49 @@ public interface Query<T> {
@Nonnull
QueryIterator<T> findIterate();
/**
* Execute the query returning the result as a Stream.
* <p>
* Note that this will hold all resulting beans in memory using a single
* persistence context. Use findLargeStream() for queries that expect to
* return a large number of results.
* </p>
* <pre>{@code
*
* // use try with resources to ensure Stream is closed
*
* try (Stream<Customer> stream = query.findStream()) {
* stream
* .map(...)
* .collect(...);
* }
*
* }</pre>
*/
@Nonnull
Stream<T> findStream();
/**
* Execute the query returning the result as a Stream.
* <p>
* Note that this uses multiple persistence contexts such that we can use
* it with a large number of results.
* </p>
* <pre>{@code
*
* // use try with resources to ensure Stream is closed
*
* try (Stream<Customer> stream = query.findLargeStream()) {
* stream
* .map(...)
* .collect(...);
* }
*
* }</pre>
*/
@Nonnull
Stream<T> findLargeStream();
/**
* Execute the query processing the beans one at a time.
* <p>
@@ -98,11 +98,11 @@ import io.ebeaninternal.server.query.CQueryEngine;
import io.ebeaninternal.server.query.CallableQueryCount;
import io.ebeaninternal.server.query.CallableQueryIds;
import io.ebeaninternal.server.query.CallableQueryList;
import io.ebeaninternal.server.query.DtoQueryEngine;
import io.ebeaninternal.server.query.LimitOffsetPagedList;
import io.ebeaninternal.server.query.QueryFutureIds;
import io.ebeaninternal.server.query.QueryFutureList;
import io.ebeaninternal.server.query.QueryFutureRowCount;
import io.ebeaninternal.server.query.DtoQueryEngine;
import io.ebeaninternal.server.querydefn.DefaultDtoQuery;
import io.ebeaninternal.server.querydefn.DefaultOrmQuery;
import io.ebeaninternal.server.querydefn.DefaultOrmUpdate;
@@ -136,11 +136,16 @@ import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.Spliterator;
import java.util.concurrent.Callable;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Consumer;
import java.util.function.Function;
import java.util.function.Predicate;
import java.util.stream.Stream;
import static java.util.Spliterators.spliteratorUnknownSize;
import static java.util.stream.StreamSupport.stream;
/**
* The default server side implementation of EbeanServer.
@@ -651,7 +656,7 @@ public final class DefaultServer implements SpiServer, SpiEbeanServer {
executeSql(connection, databasePlatform.truncateStatement(table));
}
connection.commit();
} catch(SQLException e) {
} catch (SQLException e) {
throw new PersistenceException("Error executing truncate", e);
}
}
@@ -1507,6 +1512,48 @@ public final class DefaultServer implements SpiServer, SpiEbeanServer {
}
}
@Override
public <T> Stream<T> findLargeStream(Query<T> query, Transaction transaction) {
return findStreamWithSingleContext(false, query, transaction);
}
@Override
public <T> Stream<T> findStream(Query<T> query, Transaction transaction) {
return findStreamWithSingleContext(true, query, transaction);
}
private <T> Stream<T> findStreamWithSingleContext(boolean singleContext, Query<T> query, Transaction transaction) {
SpiOrmQueryRequest<T> request = createQueryRequest(Type.ITERATE, query, transaction);
if (singleContext) {
request.setIterateSingleContext();
}
try {
request.initTransIfRequired();
return toStream(request.findIterate());
} catch (RuntimeException ex) {
request.endTransIfRequired();
throw ex;
}
}
private <T> Stream<T> toStream(QueryIterator<T> queryIterator) {
return stream(spliteratorUnknownSize(queryIterator, Spliterator.ORDERED), false)
.onClose(new QueryIteratorClose(queryIterator));
}
private static class QueryIteratorClose implements Runnable {
private final QueryIterator<?> iterator;
private QueryIteratorClose(QueryIterator<?> iterator) {
this.iterator = iterator;
}
@Override
public void run() {
iterator.close();
}
}
@Override
public <T> void findEach(Query<T> query, Consumer<T> consumer, Transaction t) {
@@ -87,6 +87,8 @@ public final class OrmQueryRequest<T> extends BeanRequest implements SpiOrmQuery
private Set<String> dependentTables;
private boolean iterateSingleContext;
/**
* Create the InternalQueryRequest.
*/
@@ -104,6 +106,15 @@ public final class OrmQueryRequest<T> extends BeanRequest implements SpiOrmQuery
return queryEngine.translate(this, bindLog, sql, e);
}
@Override
public void setIterateSingleContext() {
this.iterateSingleContext = true;
}
public boolean isIterateSingleContext() {
return iterateSingleContext;
}
@Override
public boolean isDeleteByStatement() {
if (!transaction.isPersistCascade() || beanDescriptor.isDeleteByStatement()) {
@@ -322,11 +333,13 @@ public final class OrmQueryRequest<T> extends BeanRequest implements SpiOrmQuery
* For iterate queries reset the persistenceContext and loadContext.
*/
public void flushPersistenceContextOnIterate() {
persistenceContext = new DefaultPersistenceContext();
loadContext.resetPersistenceContext(persistenceContext);
if (jsonRead != null) {
jsonRead.setPersistenceContext(persistenceContext);
jsonRead.setLoadContext(loadContext);
if (!iterateSingleContext) {
persistenceContext = new DefaultPersistenceContext();
loadContext.resetPersistenceContext(persistenceContext);
if (jsonRead != null) {
jsonRead.setPersistenceContext(persistenceContext);
jsonRead.setLoadContext(loadContext);
}
}
}
@@ -171,4 +171,10 @@ public interface SpiOrmQueryRequest<T> extends BeanQueryRequest<T>, DocQueryRequ
* Return true if delete by statement is allowed for this type given cascade rules etc.
*/
boolean isDeleteByStatement();
/**
* Set when we want to use a single persistence context for all beans returned
* in the query (so all beans are held in memory)
*/
void setIterateSingleContext();
}
@@ -204,12 +204,14 @@ public class CQueryEngine {
CQuery<T> cquery = queryBuilder.buildQuery(request);
request.setCancelableQuery(cquery);
try {
if (defaultFetchSizeFindEach > 0) {
request.setDefaultFetchBuffer(defaultFetchSizeFindEach);
}
if (!cquery.prepareBindExecuteQueryForwardOnly(forwardOnlyHintOnFindIterate)) {
if (request.isIterateSingleContext()) {
// expected relatively small number of results, single persistence context
cquery.prepareBindExecuteQuery();
} else if (!cquery.prepareBindExecuteQueryForwardOnly(forwardOnlyHintOnFindIterate)) {
// query has been cancelled already
logger.trace("Future fetch already cancelled");
return null;
@@ -37,6 +37,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.function.Consumer;
import java.util.function.Predicate;
import java.util.stream.Stream;
/**
* Implementation of FetchGroup query for use to create FetchGroup via query beans.
@@ -243,6 +244,18 @@ class DefaultFetchGroupQuery<T> implements SpiFetchGroupQuery<T> {
throw new RuntimeException("EB102: Only select() and fetch() clause is allowed on FetchGroup");
}
@Nonnull
@Override
public Stream<T> findStream() {
throw new RuntimeException("EB102: Only select() and fetch() clause is allowed on FetchGroup");
}
@Nonnull
@Override
public Stream<T> findLargeStream() {
throw new RuntimeException("EB102: Only select() and fetch() clause is allowed on FetchGroup");
}
@Override
public void findEach(Consumer<T> consumer) {
throw new RuntimeException("EB102: Only select() and fetch() clause is allowed on FetchGroup");
@@ -104,9 +104,7 @@ public class DefaultOrmQueryEngine implements OrmQueryEngine {
@Override
public <T> QueryIterator<T> findIterate(OrmQueryRequest<T> request) {
// LIMITATION: You can not use QueryIterator to load bean cache
flushJdbcBatchOnQuery(request);
return queryEngine.findIterate(request);
}
@@ -71,6 +71,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.function.Consumer;
import java.util.function.Predicate;
import java.util.stream.Stream;
/**
* Default implementation of an Object Relational query.
@@ -756,7 +757,7 @@ public class DefaultOrmQuery<T> implements SpiQuery<T> {
if (underlyingList.size() == 1) {
SpiExpression singleExpression = underlyingList.get(0);
if (singleExpression instanceof IdInExpression) {
return new CacheIdLookup<>((IdInExpression)singleExpression);
return new CacheIdLookup<>((IdInExpression) singleExpression);
}
}
return null;
@@ -1455,7 +1456,7 @@ public class DefaultOrmQuery<T> implements SpiQuery<T> {
@Override
public Query<T> usingTransaction(Transaction transaction) {
this.transaction = (SpiTransaction)transaction;
this.transaction = (SpiTransaction) transaction;
return this;
}
@@ -1521,6 +1522,16 @@ public class DefaultOrmQuery<T> implements SpiQuery<T> {
return server.findIterate(this, transaction);
}
@Override
public Stream<T> findStream() {
return server.findStream(this, transaction);
}
@Override
public Stream<T> findLargeStream() {
return server.findLargeStream(this, transaction);
}
@Override
public List<Version<T>> findVersions() {
this.temporalMode = TemporalMode.VERSIONS;
@@ -66,6 +66,7 @@ import java.util.Set;
import java.util.concurrent.Callable;
import java.util.function.Consumer;
import java.util.function.Predicate;
import java.util.stream.Stream;
/**
@@ -659,6 +660,16 @@ public class TDSpiEbeanServer implements SpiEbeanServer {
return null;
}
@Override
public <T> Stream<T> findStream(Query<T> query, Transaction transaction) {
return null;
}
@Override
public <T> Stream<T> findLargeStream(Query<T> query, Transaction transaction) {
return null;
}
@Override
public <T> void findEach(Query<T> query, Consumer<T> consumer, Transaction transaction) {
}
@@ -0,0 +1,71 @@
package org.tests.query;
import io.ebean.BaseTestCase;
import io.ebean.DB;
import org.junit.Test;
import org.tests.model.basic.Customer;
import org.tests.model.basic.ResetBasicData;
import java.util.List;
import java.util.stream.Stream;
import static java.util.stream.Collectors.toList;
import static org.assertj.core.api.Assertions.assertThat;
public class TestQueryFindStream extends BaseTestCase {
@Test
public void findStream_basic() {
ResetBasicData.reset();
try (Stream<Customer> stream = DB.find(Customer.class)
.findStream()) {
// bad example, don't use a stream like this when we can
// use findSingleAttributeList() instead
final List<String> namesStream = stream
.map(Customer::getName)
.collect(toList());
final List<String> namesQuery = DB.find(Customer.class)
.select("name")
.findSingleAttributeList();
assertThat(namesStream).hasSize(namesQuery.size());
assertThat(namesStream).containsAll(namesQuery);
}
}
@Test
public void findLargeStream_basic() {
ResetBasicData.reset();
try (Stream<Customer> stream = DB.find(Customer.class)
.findLargeStream()) {
// bad example, don't use a stream like this when we can
// use findSingleAttributeList() instead
final List<String> namesStream = stream
.map(Customer::getName)
.collect(toList());
final List<String> namesQuery = DB.find(Customer.class)
.select("name")
.findSingleAttributeList();
assertThat(namesStream).hasSize(namesQuery.size());
assertThat(namesStream).containsAll(namesQuery);
}
}
@Test
public void manualTest_findSteam_when_streamNotClosed_connectionLeak() {
Stream<Customer> stream = DB.find(Customer.class).findStream();
// remember a steam MUST be closed or we leak resources
// comment out the close(); below to leak a connection
stream.close();
}
}
+1
View File
@@ -23,6 +23,7 @@ ebean.ddl.run=true
ebean.ddl.header=-- Generated by ebean ${version} at ${timestamp}
ebean.packages=org.tests
datasource.default=h2
#datasource.h2.capturestacktrace=true
ebean.dumpMetricsOnShutdown=true
ebean.dumpMetricsOptions=sql,hash