diff --git a/ebean-api/src/main/java/io/ebean/SqlQuery.java b/ebean-api/src/main/java/io/ebean/SqlQuery.java index 8b5452ee7..46b26c90a 100644 --- a/ebean-api/src/main/java/io/ebean/SqlQuery.java +++ b/ebean-api/src/main/java/io/ebean/SqlQuery.java @@ -365,5 +365,10 @@ public interface SqlQuery extends Serializable { * Return the list of values. */ List findList(); + + /** + * Find streaming the result effectively consuming a row at a time. + */ + void findEach(Consumer consumer); } } diff --git a/ebean-core/src/main/java/io/ebeaninternal/api/SpiEbeanServer.java b/ebean-core/src/main/java/io/ebeaninternal/api/SpiEbeanServer.java index deb9c281c..a2e572904 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/api/SpiEbeanServer.java +++ b/ebean-core/src/main/java/io/ebeaninternal/api/SpiEbeanServer.java @@ -234,6 +234,11 @@ public interface SpiEbeanServer extends ExtendedServer, EbeanServer, BeanCollect */ List findSingleAttributeList(SpiSqlQuery query, Class cls); + /** + * SqlQuery find single attribute streaming the result to a consumer. + */ + void findSingleAttributeEach(SpiSqlQuery query, Class cls, Consumer consumer); + /** * SqlQuery find one with mapper. */ diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/core/DefaultServer.java b/ebean-core/src/main/java/io/ebeaninternal/server/core/DefaultServer.java index 6595d9816..81c16cf41 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/core/DefaultServer.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/core/DefaultServer.java @@ -1606,6 +1606,11 @@ public final class DefaultServer implements SpiServer, SpiEbeanServer { return executeSqlQuery((req) -> req.findOneMapper(mapper), query); } + @Override + public void findSingleAttributeEach(SpiSqlQuery query, Class cls, Consumer consumer) { + executeSqlQuery((req) -> req.findSingleAttributeEach(cls, consumer), query); + } + @Override public List findSingleAttributeList(SpiSqlQuery query, Class cls) { return executeSqlQuery((req) -> req.findSingleAttributeList(cls), query); diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryEngine.java b/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryEngine.java index 4e34a3894..e9d18797a 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryEngine.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryEngine.java @@ -30,12 +30,12 @@ public interface RelationalQueryEngine { /** * Find each via raw consumer. */ - void findEachRow(RelationalQueryRequest request, RowConsumer mapper); + void findEach(RelationalQueryRequest request, RowConsumer mapper); /** * Find one via mapper. */ - T findOneMapper(RelationalQueryRequest request, RowMapper mapper); + T findOne(RelationalQueryRequest request, RowMapper mapper); /** * Find single attribute. @@ -47,6 +47,11 @@ public interface RelationalQueryEngine { */ List findSingleAttributeList(RelationalQueryRequest request, Class cls); + /** + * Find single attribute streaming the result to a consumer. + */ + void findSingleAttributeEach(RelationalQueryRequest request, Class cls, Consumer consumer); + /** * Collect SQL query execution statistics. */ diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java b/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java index 250186c2d..9a0ceb153 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java @@ -53,7 +53,7 @@ public final class RelationalQueryRequest extends AbstractSqlQueryRequest { boolean findEachRow(RowConsumer mapper) { flushJdbcBatchOnQuery(); - queryEngine.findEachRow(this, mapper); + queryEngine.findEach(this, mapper); return true; } @@ -64,7 +64,13 @@ public final class RelationalQueryRequest extends AbstractSqlQueryRequest { T findOneMapper(RowMapper mapper) { flushJdbcBatchOnQuery(); - return queryEngine.findOneMapper(this, mapper); + return queryEngine.findOne(this, mapper); + } + + public boolean findSingleAttributeEach(Class cls, Consumer consumer) { + flushJdbcBatchOnQuery(); + queryEngine.findSingleAttributeEach(this, cls, consumer); + return true; } public List findSingleAttributeList(Class cls) { @@ -79,7 +85,7 @@ public final class RelationalQueryRequest extends AbstractSqlQueryRequest { public void findEach(Consumer consumer) { flushJdbcBatchOnQuery(); - queryEngine.findEachRow(this, (resultSet, rowNum) -> consumer.accept(createNewRow())); + queryEngine.findEach(this, (resultSet, rowNum) -> consumer.accept(createNewRow())); } public void findEachWhile(Predicate consumer) { diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/query/DefaultRelationalQueryEngine.java b/ebean-core/src/main/java/io/ebeaninternal/server/query/DefaultRelationalQueryEngine.java index b150c9c4b..78b9717d9 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/query/DefaultRelationalQueryEngine.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/query/DefaultRelationalQueryEngine.java @@ -17,6 +17,7 @@ import io.ebeaninternal.server.persist.Binder; import javax.persistence.PersistenceException; import java.util.ArrayList; import java.util.List; +import java.util.function.Consumer; import java.util.function.Predicate; /** @@ -59,7 +60,7 @@ public class DefaultRelationalQueryEngine implements RelationalQueryEngine { } @Override - public void findEachRow(RelationalQueryRequest request, RowConsumer consumer) { + public void findEach(RelationalQueryRequest request, RowConsumer consumer) { try { request.executeSql(binder, SpiQuery.Type.ITERATE); request.mapEach(consumer); @@ -93,7 +94,7 @@ public class DefaultRelationalQueryEngine implements RelationalQueryEngine { } @Override - public T findOneMapper(RelationalQueryRequest request, RowMapper mapper) { + public T findOne(RelationalQueryRequest request, RowMapper mapper) { try { request.executeSql(binder, SpiQuery.Type.BEAN); T value = request.mapOne(mapper); @@ -170,4 +171,23 @@ public class DefaultRelationalQueryEngine implements RelationalQueryEngine { } } + @SuppressWarnings("unchecked") + @Override + public void findSingleAttributeEach(RelationalQueryRequest request, Class cls, Consumer consumer) { + ScalarType scalarType = (ScalarType) binder.getScalarType(cls); + try { + request.executeSql(binder, SpiQuery.Type.ATTRIBUTE); + final DataReader dataReader = binder.createDataReader(request.getResultSet()); + while (dataReader.next()) { + consumer.accept(scalarType.read(dataReader)); + } + request.logSummary(); + + } catch (Exception e) { + throw new PersistenceException(errMsg(e.getMessage(), request.getSql()), e); + + } finally { + request.close(); + } + } } diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java index 4cf446d0f..182b1cfd4 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java @@ -254,6 +254,15 @@ public class DefaultRelationalQuery implements SpiSqlQuery { public List findList() { return findSingleAttributeList(type); } + + @Override + public void findEach(Consumer consumer) { + scalarFindEach(type, consumer); + } + } + + private void scalarFindEach(Class type, Consumer consumer) { + server.findSingleAttributeEach(this, type, consumer); } private class Mapper implements SqlQuery.TypeQuery { @@ -279,7 +288,7 @@ public class DefaultRelationalQuery implements SpiSqlQuery { return mapperFindList(mapper); } - //@Override + @Override public void findEach(Consumer consumer) { mapperFindEach(mapper, consumer); } diff --git a/ebean-core/src/test/java/io/ebeaninternal/api/TDSpiEbeanServer.java b/ebean-core/src/test/java/io/ebeaninternal/api/TDSpiEbeanServer.java index 03ed3362f..71e73b59f 100644 --- a/ebean-core/src/test/java/io/ebeaninternal/api/TDSpiEbeanServer.java +++ b/ebean-core/src/test/java/io/ebeaninternal/api/TDSpiEbeanServer.java @@ -493,6 +493,10 @@ public class TDSpiEbeanServer implements SpiEbeanServer { return null; } + @Override + public void findSingleAttributeEach(SpiSqlQuery query, Class cls, Consumer consumer) { + } + @Override public T findOneMapper(SpiSqlQuery query, RowMapper mapper) { return null; diff --git a/ebean-core/src/test/java/org/tests/query/sqlquery/SqlQueryTests.java b/ebean-core/src/test/java/org/tests/query/sqlquery/SqlQueryTests.java index 7745b4ca6..572ebf39e 100644 --- a/ebean-core/src/test/java/org/tests/query/sqlquery/SqlQueryTests.java +++ b/ebean-core/src/test/java/org/tests/query/sqlquery/SqlQueryTests.java @@ -39,6 +39,28 @@ public class SqlQueryTests extends BaseTestCase { assertThat(lineAmounts).isNotEmpty(); } + @Test + public void findSingleAttributeEach_decimal() { + + ResetBasicData.reset(); + + String sql = "select (unit_price * order_qty) from o_order_detail where unit_price > ? order by (unit_price * order_qty) desc"; + + AtomicLong counter = new AtomicLong(); + AtomicLong inc = new AtomicLong(); + + DB.sqlQuery(sql) + .setParameter(3) + .mapToScalar(BigDecimal.class) + .findEach(val -> { + counter.incrementAndGet(); + inc.addAndGet(val.longValue()); + }); + + assertThat(inc.get()).isGreaterThan(counter.get()); + assertThat(counter.get()).isGreaterThan(0); + } + @Test public void findSingleDecimal() { @@ -141,6 +163,24 @@ public class SqlQueryTests extends BaseTestCase { private static final CustMapper CUST_MAPPER = new CustMapper(); + @Test + public void findEach_mapper() { + + ResetBasicData.reset(); + + String sql = "select id, name, status from o_customer where name is not null"; + + AtomicInteger counter = new AtomicInteger(); + DB.sqlQuery(sql) + .mapTo(CUST_MAPPER) + .findEach(custDto -> { + counter.incrementAndGet(); + assertThat(custDto.name).isNotNull(); + }); + + assertThat(counter.get()).isGreaterThan(0); + } + @Test public void findOne_mapper() {