From 7cbae758fcefe98b0aa27c88981841deecc9625b Mon Sep 17 00:00:00 2001 From: rob bygrave Date: Fri, 8 Jun 2018 02:04:38 +1200 Subject: [PATCH] #1407 - ENH: Add SqlQuery RowMapper and RowConsumer ... for raw JDBC ResultSet mapping and consuming --- src/main/java/io/ebean/RowConsumer.java | 40 ++++++++ src/main/java/io/ebean/RowMapper.java | 62 ++++++++++++ src/main/java/io/ebean/SqlQuery.java | 44 +++++++++ .../io/ebeaninternal/api/SpiEbeanServer.java | 17 ++++ .../server/core/DefaultServer.java | 38 +++++--- .../server/core/RelationalQueryEngine.java | 17 ++++ .../server/core/RelationalQueryRequest.java | 39 ++++++++ .../query/DefaultRelationalQueryEngine.java | 49 ++++++++++ .../querydefn/DefaultRelationalQuery.java | 17 ++++ .../ebeaninternal/api/TDSpiEbeanServer.java | 17 ++++ .../tests/query/sqlquery/SqlQueryTests.java | 94 +++++++++++++++++++ 11 files changed, 421 insertions(+), 13 deletions(-) create mode 100644 src/main/java/io/ebean/RowConsumer.java create mode 100644 src/main/java/io/ebean/RowMapper.java diff --git a/src/main/java/io/ebean/RowConsumer.java b/src/main/java/io/ebean/RowConsumer.java new file mode 100644 index 000000000..4845f18ea --- /dev/null +++ b/src/main/java/io/ebean/RowConsumer.java @@ -0,0 +1,40 @@ +package io.ebean; + +import java.sql.ResultSet; +import java.sql.SQLException; + +/** + * Used with SqlQuery to process potentially large queries reading directly from the JDBC ResultSet. + *

+ * This provides a low level option that reads directly from the JDBC ResultSet. + *

+ * + *
{@code
+ *
+ *  String sql = "select id, name, status from o_customer order by name desc";
+ *
+ *  Ebean.createSqlQuery(sql)
+ *    .findEachRow((resultSet, rowNum) -> {
+ *
+ *      // read directly from ResultSet
+ *
+ *      long id = resultSet.getLong(1);
+ *      String name = resultSet.getString(2);
+ *
+ *      // do something interesting with the data
+ *
+ *    });
+ *
+ * }
+ */ +@FunctionalInterface +public interface RowConsumer { + + /** + * Read the data from the ResultSet and process it. + * + * @param resultSet The JDBC ResultSet positioned to the current row + * @param rowNum The number of the current row being mapped. + */ + void accept(ResultSet resultSet, int rowNum) throws SQLException; +} diff --git a/src/main/java/io/ebean/RowMapper.java b/src/main/java/io/ebean/RowMapper.java new file mode 100644 index 000000000..1cc067372 --- /dev/null +++ b/src/main/java/io/ebean/RowMapper.java @@ -0,0 +1,62 @@ +package io.ebean; + +import java.sql.ResultSet; +import java.sql.SQLException; + +/** + * Used with SqlQuery to map raw JDBC ResultSet to objects. + *

+ * This provides a low level mapping option with direct use of JDBC ResultSet + * with the option of having logic in the mapping. For example, only map some + * columns depending on the values read from other columns. + *

+ *

+ * For straight mapping into beans then DtoQuery would be the first choice as + * it can automatically map the ResultSet into beans. + *

+ * + *
{@code
+ *
+ *    //
+ *    // A mapper from ResultSet into our CustomerDto bean
+ *    //
+ *    class CustomerMapper implements RowMapper {
+ *
+ *     @Override
+ *     public CustomerDto map(ResultSet rset, int rowNum) throws SQLException {
+ *
+ *       long id = rset.getLong(1);
+ *       String name = rset.getString(2);
+ *       String status = rset.getString(3);
+ *
+ *       return new CustomerDto(id, name, status);
+ *     }
+ *   }
+ *
+ *
+ *   //
+ *   // Then use the mapper
+ *   //
+ *
+ *   String sql = "select id, name, status from o_customer where name = ?";
+ *
+ *  CustomerDto rob = Ebean.createSqlQuery(sql)
+ *    .setParameter(1, "Rob")
+ *    .findOne(CUSTOMER_MAPPER);
+ *
+ *
+ * }
+ * + * @param The type the row data is mapped into. + */ +@FunctionalInterface +public interface RowMapper { + + /** + * Read the data from the ResultSet and map to the return type. + * + * @param resultSet The JDBC ResultSet positioned to the current row + * @param rowNum The number of the current row being mapped. + */ + T map(ResultSet resultSet, int rowNum) throws SQLException; +} diff --git a/src/main/java/io/ebean/SqlQuery.java b/src/main/java/io/ebean/SqlQuery.java index 8f8b9ed01..2adbd67b2 100644 --- a/src/main/java/io/ebean/SqlQuery.java +++ b/src/main/java/io/ebean/SqlQuery.java @@ -77,6 +77,50 @@ public interface SqlQuery extends Serializable { @Nullable SqlRow findOne(); + /** + * Execute the query returning a single result using the mapper. + * + * @param mapper Used to map each ResultSet row into the result object. + */ + T findOne(RowMapper mapper); + + /** + * Execute the query returning a list using the mapper. + * + * @param mapper Used to map each ResultSet row into the result object. + */ + List findList(RowMapper mapper); + + /** + * Execute the query reading each row from ResultSet using the RowConsumer. + *

+ * This provides a low level option that reads directly from the JDBC ResultSet + * and is good for processing very large results where (unlike findList) we don't + * hold all the results in memory but instead can process row by row. + *

+ * + *
{@code
+   *
+   *  String sql = "select id, name, status from customer order by name desc";
+   *
+   *  Ebean.createSqlQuery(sql)
+   *    .findEachRow((resultSet, rowNum) -> {
+   *
+   *      // read directly from ResultSet
+   *
+   *      long id = resultSet.getLong(1);
+   *      String name = resultSet.getString(2);
+   *
+   *      // do something interesting with the data
+   *
+   *    });
+   *
+   * }
+ * + * @param consumer Used to read and process each ResultSet row. + */ + void findEachRow(RowConsumer consumer); + /** * Execute the query returning an optional row. */ diff --git a/src/main/java/io/ebeaninternal/api/SpiEbeanServer.java b/src/main/java/io/ebeaninternal/api/SpiEbeanServer.java index ed6fac793..af0154a8e 100644 --- a/src/main/java/io/ebeaninternal/api/SpiEbeanServer.java +++ b/src/main/java/io/ebeaninternal/api/SpiEbeanServer.java @@ -4,6 +4,8 @@ import io.ebean.DtoQuery; import io.ebean.EbeanServer; import io.ebean.PersistenceContextScope; import io.ebean.Query; +import io.ebean.RowConsumer; +import io.ebean.RowMapper; import io.ebean.Transaction; import io.ebean.TxScope; import io.ebean.bean.BeanCollectionLoader; @@ -244,6 +246,21 @@ public interface SpiEbeanServer extends EbeanServer, BeanLoader, BeanCollectionL */ List findSingleAttributeList(SpiSqlQuery query, Class cls); + /** + * SqlQuery find one with mapper. + */ + T findOneMapper(SpiSqlQuery query, RowMapper mapper); + + /** + * SqlQuery find list with mapper. + */ + List findListMapper(SpiSqlQuery query, RowMapper mapper); + + /** + * SqlQuery find each with consumer. + */ + void findEachRow(SpiSqlQuery query, RowConsumer consumer); + /** * DTO findList query. */ diff --git a/src/main/java/io/ebeaninternal/server/core/DefaultServer.java b/src/main/java/io/ebeaninternal/server/core/DefaultServer.java index 14bb12e0f..779cb3982 100644 --- a/src/main/java/io/ebeaninternal/server/core/DefaultServer.java +++ b/src/main/java/io/ebeaninternal/server/core/DefaultServer.java @@ -19,6 +19,8 @@ import io.ebean.PersistenceContextScope; import io.ebean.ProfileLocation; import io.ebean.Query; import io.ebean.QueryIterator; +import io.ebean.RowConsumer; +import io.ebean.RowMapper; import io.ebean.SqlQuery; import io.ebean.SqlRow; import io.ebean.SqlUpdate; @@ -1586,29 +1588,39 @@ public final class DefaultServer implements SpiServer, SpiEbeanServer { } } - @Override - public List findSingleAttributeList(SpiSqlQuery query, Class cls) { + private

P executeSqlQuery(Function fun, SpiSqlQuery query) { RelationalQueryRequest request = new RelationalQueryRequest(this, relationalQueryEngine, query, null); try { request.initTransIfRequired(); - return request.findSingleAttributeList(cls); - + return fun.apply(request); } finally { request.endTransIfRequired(); } } + @Override + public void findEachRow(SpiSqlQuery query, RowConsumer consumer) { + executeSqlQuery((req) -> req.findEachRow(consumer), query); + } + + @Override + public List findListMapper(SpiSqlQuery query, RowMapper mapper) { + return executeSqlQuery((req) -> req.findListMapper(mapper), query); + } + + @Override + public T findOneMapper(SpiSqlQuery query, RowMapper mapper) { + return executeSqlQuery((req) -> req.findOneMapper(mapper), query); + } + + @Override + public List findSingleAttributeList(SpiSqlQuery query, Class cls) { + return executeSqlQuery((req) -> req.findSingleAttributeList(cls), query); + } + @Override public T findSingleAttribute(SpiSqlQuery query, Class cls) { - - RelationalQueryRequest request = new RelationalQueryRequest(this, relationalQueryEngine, query, null); - try { - request.initTransIfRequired(); - return request.findSingleAttribute(cls); - - } finally { - request.endTransIfRequired(); - } + return executeSqlQuery((req) -> req.findSingleAttribute(cls), query); } @Override diff --git a/src/main/java/io/ebeaninternal/server/core/RelationalQueryEngine.java b/src/main/java/io/ebeaninternal/server/core/RelationalQueryEngine.java index f6405aaa8..718a2873a 100644 --- a/src/main/java/io/ebeaninternal/server/core/RelationalQueryEngine.java +++ b/src/main/java/io/ebeaninternal/server/core/RelationalQueryEngine.java @@ -1,6 +1,8 @@ package io.ebeaninternal.server.core; +import io.ebean.RowConsumer; +import io.ebean.RowMapper; import io.ebean.SqlRow; import io.ebean.meta.MetricVisitor; @@ -40,6 +42,21 @@ public interface RelationalQueryEngine { */ List findSingleAttributeList(RelationalQueryRequest request, Class cls); + /** + * Find one via mapper. + */ + T findOneMapper(RelationalQueryRequest request, RowMapper mapper); + + /** + * Find list via mapper. + */ + List findListMapper(RelationalQueryRequest request, RowMapper mapper); + + /** + * Find each via raw consumer. + */ + void findEachRow(RelationalQueryRequest request, RowConsumer mapper); + /** * Collect SQL query execution statistics. */ diff --git a/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java b/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java index 6324ffdfc..c9a47fb90 100644 --- a/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java +++ b/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java @@ -1,5 +1,7 @@ package io.ebeaninternal.server.core; +import io.ebean.RowConsumer; +import io.ebean.RowMapper; import io.ebean.SqlQuery; import io.ebean.SqlRow; import io.ebean.Transaction; @@ -53,6 +55,19 @@ public final class RelationalQueryRequest extends AbstractSqlQueryRequest { } } + boolean findEachRow(RowConsumer mapper) { + queryEngine.findEachRow(this, mapper); + return true; + } + + List findListMapper(RowMapper mapper) { + return queryEngine.findListMapper(this, mapper); + } + + T findOneMapper(RowMapper mapper) { + return queryEngine.findOneMapper(this, mapper); + } + public List findSingleAttributeList(Class cls) { return queryEngine.findSingleAttributeList(this, cls); } @@ -119,4 +134,28 @@ public final class RelationalQueryRequest extends AbstractSqlQueryRequest { public void incrementRows() { rows++; } + + public List mapList(RowMapper mapper) throws SQLException { + + List list = new ArrayList<>(); + while (next()) { + list.add(mapper.map(resultSet, rows++)); + } + return list; + } + + public T mapOne(RowMapper mapper) throws SQLException { + if (!next()) { + return null; + } else { + return mapper.map(resultSet, rows++); + } + } + + public void mapEach(RowConsumer consumer) throws SQLException { + while (next()) { + consumer.accept(resultSet, rows++); + } + } + } diff --git a/src/main/java/io/ebeaninternal/server/query/DefaultRelationalQueryEngine.java b/src/main/java/io/ebeaninternal/server/query/DefaultRelationalQueryEngine.java index b2fe0466e..d46aaf0a3 100644 --- a/src/main/java/io/ebeaninternal/server/query/DefaultRelationalQueryEngine.java +++ b/src/main/java/io/ebeaninternal/server/query/DefaultRelationalQueryEngine.java @@ -1,5 +1,7 @@ package io.ebeaninternal.server.query; +import io.ebean.RowConsumer; +import io.ebean.RowMapper; import io.ebean.SqlRow; import io.ebean.meta.MetricType; import io.ebean.meta.MetricVisitor; @@ -92,6 +94,53 @@ public class DefaultRelationalQueryEngine implements RelationalQueryEngine { } } + @Override + public T findOneMapper(RelationalQueryRequest request, RowMapper mapper) { + try { + request.executeSql(binder, SpiQuery.Type.BEAN); + T value = request.mapOne(mapper); + request.logSummary(); + return value; + + } catch (Exception e) { + throw new PersistenceException(Message.msg("fetch.error", e.getMessage(), request.getSql()), e); + + } finally { + request.close(); + } + } + + @Override + public List findListMapper(RelationalQueryRequest request, RowMapper mapper) { + try { + request.executeSql(binder, SpiQuery.Type.LIST); + List list = request.mapList(mapper); + request.logSummary(); + return list; + + } catch (Exception e) { + throw new PersistenceException(Message.msg("fetch.error", e.getMessage(), request.getSql()), e); + + } finally { + request.close(); + } + } + + @Override + public void findEachRow(RelationalQueryRequest request, RowConsumer consumer) { + try { + request.executeSql(binder, SpiQuery.Type.LIST); + request.mapEach(consumer); + request.logSummary(); + + } catch (Exception e) { + throw new PersistenceException(Message.msg("fetch.error", e.getMessage(), request.getSql()), e); + + } finally { + request.close(); + } + } + @SuppressWarnings("unchecked") @Override public List findSingleAttributeList(RelationalQueryRequest request, Class cls) { diff --git a/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java b/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java index ef58e2804..5ca2e1075 100644 --- a/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java +++ b/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java @@ -1,5 +1,7 @@ package io.ebeaninternal.server.querydefn; +import io.ebean.RowConsumer; +import io.ebean.RowMapper; import io.ebean.SqlRow; import io.ebeaninternal.api.BindParams; import io.ebeaninternal.api.SpiEbeanServer; @@ -69,6 +71,21 @@ public class DefaultRelationalQuery implements SpiSqlQuery { return server.findSingleAttributeList(this, cls); } + @Override + public T findOne(RowMapper mapper) { + return server.findOneMapper(this, mapper); + } + + @Override + public List findList(RowMapper mapper) { + return server.findListMapper(this, mapper); + } + + @Override + public void findEachRow(RowConsumer consumer) { + server.findEachRow(this, consumer); + } + @Override public SqlRow findOne() { return server.findOne(this, null); diff --git a/src/test/java/io/ebeaninternal/api/TDSpiEbeanServer.java b/src/test/java/io/ebeaninternal/api/TDSpiEbeanServer.java index 1661403b9..64bff0d26 100644 --- a/src/test/java/io/ebeaninternal/api/TDSpiEbeanServer.java +++ b/src/test/java/io/ebeaninternal/api/TDSpiEbeanServer.java @@ -16,6 +16,8 @@ import io.ebean.PagedList; import io.ebean.PersistenceContextScope; import io.ebean.Query; import io.ebean.QueryIterator; +import io.ebean.RowConsumer; +import io.ebean.RowMapper; import io.ebean.SqlQuery; import io.ebean.SqlRow; import io.ebean.SqlUpdate; @@ -467,6 +469,21 @@ public class TDSpiEbeanServer implements SpiEbeanServer { return null; } + @Override + public T findOneMapper(SpiSqlQuery query, RowMapper mapper) { + return null; + } + + @Override + public List findListMapper(SpiSqlQuery query, RowMapper mapper) { + return null; + } + + @Override + public void findEachRow(SpiSqlQuery query, RowConsumer consumer) { + + } + @Override public SqlQuery createSqlQuery(String sql) { return null; diff --git a/src/test/java/org/tests/query/sqlquery/SqlQueryTests.java b/src/test/java/org/tests/query/sqlquery/SqlQueryTests.java index 2d4251db7..3cf065b4b 100644 --- a/src/test/java/org/tests/query/sqlquery/SqlQueryTests.java +++ b/src/test/java/org/tests/query/sqlquery/SqlQueryTests.java @@ -2,6 +2,7 @@ package org.tests.query.sqlquery; import io.ebean.BaseTestCase; import io.ebean.Ebean; +import io.ebean.RowMapper; import io.ebean.SqlQuery; import io.ebean.SqlRow; import io.ebean.meta.MetaTimedMetric; @@ -11,9 +12,12 @@ import org.tests.model.basic.Order; import org.tests.model.basic.ResetBasicData; import java.math.BigDecimal; +import java.sql.ResultSet; +import java.sql.SQLException; import java.time.OffsetDateTime; import java.util.List; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.assertEquals; @@ -76,6 +80,96 @@ public class SqlQueryTests extends BaseTestCase { assertThat(minCreated).isBefore(OffsetDateTime.now()); } + static class CustDto { + + long id; + String name; + String status; + + public CustDto(long id, String name, String status) { + this.id = id; + this.name = name; + this.status = status; + } + } + + static class CustMapper implements RowMapper { + + @Override + public CustDto map(ResultSet rset, int rowNum) throws SQLException { + + long id = rset.getLong(1); + String name = rset.getString(2); + String status = rset.getString(3); + + return new CustDto(id, name, status); + } + } + + private static final CustMapper CUST_MAPPER = new CustMapper(); + + @Test + public void findOne_mapper() { + + ResetBasicData.reset(); + + String sql = "select id, name, status from o_customer where name = ?"; + + CustDto rob = Ebean.createSqlQuery(sql) + .setParameter(1, "Rob") + .findOne(CUST_MAPPER); + + assertThat(rob.name).isEqualTo("Rob"); + } + + @Test + public void findList_mapper() { + + ResetBasicData.reset(); + + String sql = "select id, name, status from o_customer order by name desc"; + + List dtos = Ebean.createSqlQuery(sql) + .findList(CUST_MAPPER); + + assertThat(dtos).isNotEmpty(); + } + + @Test + public void findEachRow() { + + ResetBasicData.reset(); + + String sql = "select id, name, status from o_customer order by name desc"; + + AtomicLong count = new AtomicLong(); + + Ebean.createSqlQuery(sql) + .findEachRow((resultSet, rowNum) -> { + count.incrementAndGet(); + + long id = resultSet.getLong(1); + String name = resultSet.getString(2); + + System.out.println("rowNum:" + rowNum + " id:" + id + " name:" + name); + }); + + assertThat(count.get()).isGreaterThan(0); + } + + @Test + public void findOne_mapper_lambda() { + + ResetBasicData.reset(); + + String sql = "select max(id) from o_customer where name != ?"; + + long maxId = Ebean.createSqlQuery(sql) + .setParameter(1, "Rob") + .findOne((resultSet, rowNum) -> resultSet.getLong(1)); + + assertThat(maxId).isGreaterThan(0); + } @Test public void newline_replacedInLogsOnly() {