From 39928948f0a637d9512baa017d9e9736db9e8a91 Mon Sep 17 00:00:00 2001 From: John Blum Date: Fri, 20 May 2016 12:40:16 -0700 Subject: [PATCH] DATACASS-286 - Log all CQL queries executed with CqlTemplate. Additional refactoring to simplify CqlTemplate logic as well as make CqlTemplate methods more behaviorally consistent. Origina pull request: #61. --- .../cassandra/core/CqlOperations.java | 20 +- .../cassandra/core/CqlTemplate.java | 1387 ++++++++--------- .../cassandra/support/CassandraAccessor.java | 58 +- .../cassandra/core/CqlTemplateUnitTests.java | 267 ++++ .../support/CassandraAccessorUnitTests.java | 108 ++ .../core/CqlOperationsIntegrationTests.java | 8 +- .../cassandra/core/CassandraOperations.java | 2 +- .../cassandra/core/CassandraTemplate.java | 17 +- 8 files changed, 1123 insertions(+), 744 deletions(-) create mode 100644 spring-cql/src/test/java/org/springframework/cassandra/core/CqlTemplateUnitTests.java create mode 100644 spring-cql/src/test/java/org/springframework/cassandra/support/CassandraAccessorUnitTests.java diff --git a/spring-cql/src/main/java/org/springframework/cassandra/core/CqlOperations.java b/spring-cql/src/main/java/org/springframework/cassandra/core/CqlOperations.java index a936b1c00..5925f9228 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/core/CqlOperations.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/core/CqlOperations.java @@ -31,6 +31,7 @@ import org.springframework.cassandra.core.keyspace.DropIndexSpecification; import org.springframework.cassandra.core.keyspace.DropKeyspaceSpecification; import org.springframework.cassandra.core.keyspace.DropTableSpecification; import org.springframework.dao.DataAccessException; +import org.springframework.dao.IncorrectResultSizeDataAccessException; import com.datastax.driver.core.PreparedStatement; import com.datastax.driver.core.ResultSet; @@ -693,14 +694,17 @@ public interface CqlOperations { T queryForObject(Select select, RowMapper rowMapper) throws DataAccessException; /** - * Process a ResultSet through a RowMapper. This is used internal to the Template for core operations, but is made - * available through Operations in the event you have a ResultSet to process. The ResultsSet could come from a - * ResultSetFuture after an asynchronous query. - * - * @param resultSet - * @param rowMapper - * @return - * @throws DataAccessException + * Process {@link ResultSet} with {@link RowMapper}. This method is used internally to the template + * for core operations, but is made available through this interface in the event you have a {@link ResultSet} + * to process. The {@link ResultSet} could come from a {@link ResultSetFuture} after an asynchronous query. + * + * @param resultSet {@link ResultSet} to process. + * @param rowMapper {@link RowMapper} used to process the single row of the result set. + * @throws IllegalArgumentException if {@link ResultSet} is null. + * @throws IncorrectResultSizeDataAccessException if no rows are found, or more than 1 row is found. + * @throws DataAccessException if a Cassandra driver error occurs. + * @see org.springframework.cassandra.core.RowMapper + * @see com.datastax.driver.core.ResultSet */ T processOne(ResultSet resultSet, RowMapper rowMapper) throws DataAccessException; diff --git a/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java b/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java index 1cd145a67..e1796588f 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java @@ -13,9 +13,10 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.cassandra.core; -import static org.springframework.cassandra.core.cql.CqlIdentifier.cqlId; +import static org.springframework.cassandra.core.cql.CqlIdentifier.*; import java.util.ArrayList; import java.util.Collection; @@ -23,14 +24,13 @@ import java.util.HashMap; import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.NoSuchElementException; import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.concurrent.Executor; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import org.springframework.cassandra.core.cql.CqlIdentifier; import org.springframework.cassandra.core.cql.generator.AlterKeyspaceCqlGenerator; import org.springframework.cassandra.core.cql.generator.AlterTableCqlGenerator; @@ -50,6 +50,7 @@ import org.springframework.cassandra.core.keyspace.DropKeyspaceSpecification; import org.springframework.cassandra.core.keyspace.DropTableSpecification; import org.springframework.cassandra.support.CassandraAccessor; import org.springframework.dao.DataAccessException; +import org.springframework.dao.IncorrectResultSizeDataAccessException; import org.springframework.dao.InvalidDataAccessApiUsageException; import org.springframework.dao.QueryTimeoutException; import org.springframework.util.Assert; @@ -67,6 +68,7 @@ import com.datastax.driver.core.Row; import com.datastax.driver.core.Session; import com.datastax.driver.core.SimpleStatement; import com.datastax.driver.core.Statement; +import com.datastax.driver.core.TypeCodec; import com.datastax.driver.core.exceptions.DriverException; import com.datastax.driver.core.querybuilder.Batch; import com.datastax.driver.core.querybuilder.Delete; @@ -77,138 +79,176 @@ import com.datastax.driver.core.querybuilder.Truncate; import com.datastax.driver.core.querybuilder.Update; /** - * This is the Central class in the Cassandra core package. It simplifies the use of Cassandra and helps to avoid - * common errors. It executes the core Cassandra workflow, leaving application code to provide CQL and result - * extraction. This class execute CQL Queries, provides different ways to extract/map results, and provides Exception - * translation to the generic, more informative exception hierarchy defined in the org.springframework.dao - * package. + * This is the central class in the Cassandra core package. {@link CqlTemplate} simplifies the use of Cassandra + * and helps to avoid common errors. The template executes the core Cassandra workflow, leaving application code + * to provide CQL and result handling. The template executes CQL queries, provides different ways to extract and map + * results, and provides Exception translation to the generic, more informative exception hierarchy defined in the + * org.springframework.dao package. *

- * For working with POJOs, use the {@link CassandraTemplate}. + * For working with POJOs, use the CassandraTemplate. *

- * * @author David Webb * @author Matthew Adams * @author Ryan Scheidter * @author Antoine Toulme + * @author John Blum + * @see org.springframework.cassandra.core.CqlOperations + * @see org.springframework.cassandra.support.CassandraAccessor */ public class CqlTemplate extends CassandraAccessor implements CqlOperations { - protected static final Logger log = LoggerFactory.getLogger(CqlTemplate.class); - - /** - * Add common {@link QueryOptions} options for all types of queries. - * - * @param q - * @param options - * @return the {@link Statement} given. - */ - public static Statement addQueryOptions(Statement q, QueryOptions options) { - - if (options == null) { - return q; + protected static final Executor RUN_RUNNABLE_EXECUTOR = new Executor() { + @Override + @SuppressWarnings("all") + public void execute(Runnable command) { + command.run(); } + }; - if (options.getConsistencyLevel() != null) { - q.setConsistencyLevel(ConsistencyLevelResolver.resolve(options.getConsistencyLevel())); - } - if (options.getRetryPolicy() != null) { - q.setRetryPolicy(RetryPolicyResolver.resolve(options.getRetryPolicy())); - } + protected static final ResultSetExtractor RESULT_SET_RETURNING_EXTRACTOR = + new ResultSetExtractor() { + @Override + public ResultSet extractData(ResultSet resultSet) throws DriverException, DataAccessException { + return resultSet; + } + }; - return q; + /* (non-Javadoc) */ + protected String logCql(String cql) { + return logCql("executing CQL [{}]", cql); + } + + /* (non-Javadoc) */ + protected String logCql(String message, String cql) { + logDebug(message, cql); + return cql; } /** - * Add common {@link WriteOptions} options for {@link Insert} queries. + * Add common {@link QueryOptions} to all types of queries. * - * @param q - * @param options - * @return the {@link Insert} given. + * @param statement CQL {@link Statement} to execute. + * @param queryOptions query options (e.g. consistency level) to add to the CQL statement. + * @return the given {@link Statement}. */ - public static Insert addWriteOptions(Insert q, WriteOptions options) { + public static Statement addQueryOptions(Statement statement, QueryOptions queryOptions) { - if (options == null) { - return q; + if (queryOptions != null) { + if (queryOptions.getConsistencyLevel() != null) { + statement.setConsistencyLevel(ConsistencyLevelResolver.resolve(queryOptions.getConsistencyLevel())); + } + if (queryOptions.getRetryPolicy() != null) { + statement.setRetryPolicy(RetryPolicyResolver.resolve(queryOptions.getRetryPolicy())); + } } - if (options.getConsistencyLevel() != null) { - q.setConsistencyLevel(ConsistencyLevelResolver.resolve(options.getConsistencyLevel())); - } - if (options.getRetryPolicy() != null) { - q.setRetryPolicy(RetryPolicyResolver.resolve(options.getRetryPolicy())); - } - if (options.getTtl() != null) { - q.using(QueryBuilder.ttl(options.getTtl())); - } - - return q; + return statement; } /** - * Add common {@link WriteOptions} options for {@link Update} queries. + * Add common {@link WriteOptions} options to {@link Insert} CQL statements. * - * @param q - * @param options - * @return the {@link Update} given. + * @param insert {@link Insert} CQL statement to execute. + * @param writeOptions write options (e.g. consistency level) to add to the CQL statement. + * @return the given {@link Insert}. */ - public static Update addWriteOptions(Update q, WriteOptions options) { + public static Insert addWriteOptions(Insert insert, WriteOptions writeOptions) { - if (options == null) { - return q; + if (writeOptions != null) { + addQueryOptions(insert, writeOptions); + + if (writeOptions.getTtl() != null) { + insert.using(QueryBuilder.ttl(writeOptions.getTtl())); + } } - if (options.getConsistencyLevel() != null) { - q.setConsistencyLevel(ConsistencyLevelResolver.resolve(options.getConsistencyLevel())); - } - if (options.getRetryPolicy() != null) { - q.setRetryPolicy(RetryPolicyResolver.resolve(options.getRetryPolicy())); - } - if (options.getTtl() != null) { - q.using(QueryBuilder.ttl(options.getTtl())); - } - - return q; + return insert; } /** - * Add common Query options for all types of queries. + * Add common {@link WriteOptions} options to {@link Update} CQL statements. * - * @param s the prepared statement - * @param options + * @param update {@link Update} CQL statement to execute. + * @param writeOptions write options (e.g. consistency level) to add to the CQL statement. + * @return the given {@link Update}. */ - public static void addPreparedStatementOptions(PreparedStatement s, QueryOptions options) { + public static Update addWriteOptions(Update update, WriteOptions writeOptions) { - if (options == null) { - return; - } - - /* - * Add Query Options - */ - if (options.getConsistencyLevel() != null) { - s.setConsistencyLevel(ConsistencyLevelResolver.resolve(options.getConsistencyLevel())); - } - if (options.getRetryPolicy() != null) { - s.setRetryPolicy(RetryPolicyResolver.resolve(options.getRetryPolicy())); + if (writeOptions != null) { + addQueryOptions(update, writeOptions); + + if (writeOptions.getTtl() != null) { + update.using(QueryBuilder.ttl(writeOptions.getTtl())); + } } + return update; } /** - * Blank constructor. You must wire in the Session before use. + * Add common {@link QueryOptions} to Cassandra {@link PreparedStatement}s. + * + * @param preparedStatement the Cassandra {@link PreparedStatement} to execute. + * @param queryOptions query options (e.g. consistency level) to add to the Cassandra {@link PreparedStatement}. + */ + public static PreparedStatement addPreparedStatementOptions(PreparedStatement preparedStatement, + QueryOptions queryOptions) { + + if (queryOptions != null) { + if (queryOptions.getConsistencyLevel() != null) { + preparedStatement.setConsistencyLevel(ConsistencyLevelResolver.resolve( + queryOptions.getConsistencyLevel())); + } + if (queryOptions.getRetryPolicy() != null) { + preparedStatement.setRetryPolicy(RetryPolicyResolver.resolve( + queryOptions.getRetryPolicy())); + } + } + + return preparedStatement; + } + + /** + * Constructs an uninitialized instance of {@link CqlTemplate}. A Cassandra {@link Session} + * is required before use. + * + * @see #CqlTemplate(Session) */ public CqlTemplate() { } /** - * Constructor used for a basic template configuration + * Constructs an instance of {@link CqlTemplate} initialized with the given {@link Session}. * - * @param session must not be {@literal null}. + * @param session Cassandra {@link Session} used by this template to perform CQL operations. + * Must not be {@literal null}. + * @see com.datastax.driver.core.Session + * @see #setSession(Session) */ + // TODO should probably not call setSession(..) in constructor for initialization safety; + // only really matters if CqlTemplate makes Thread-safety guarantees, which currently it does not. public CqlTemplate(Session session) { setSession(session); } + /** + * Executes the given command in a Cassandra {@link Session}. + * + * @param Class type of the callback return value. + * @param callback {@link SessionCallback} to execute in the context of a Cassandra {@link Session}. + * @return the result of the callback. + */ + protected T doExecute(SessionCallback callback) { + + Assert.notNull(callback); + + try { + return callback.doInSession(getSession()); + } catch (DataAccessException e) { + throw translateExceptionIfPossible(e); + } + } + @Override public T execute(SessionCallback sessionCallback) throws DataAccessException { return doExecute(sessionCallback); @@ -225,39 +265,44 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } @Override - public void execute(Statement query) throws DataAccessException { - doExecute(query); + public void execute(Statement statement) throws DataAccessException { + doExecute(statement); } @Override public ResultSetFuture queryAsynchronously(final String cql) { return execute(new SessionCallback() { @Override - public ResultSetFuture doInSession(Session s) throws DataAccessException { - return s.executeAsync(cql); + public ResultSetFuture doInSession(Session session) throws DataAccessException { + return session.executeAsync(logCql("async execute CQL [{}]", cql)); } }); } @Override - public T queryAsynchronously(String cql, ResultSetExtractor rse, Long timeout, TimeUnit timeUnit) { - return queryAsynchronously(cql, rse, timeout, timeUnit, null); + public T queryAsynchronously(String cql, ResultSetExtractor resultSetExtractor, + Long timeout, TimeUnit timeUnit) { + + return queryAsynchronously(cql, resultSetExtractor, timeout, timeUnit, null); } @Override - public T queryAsynchronously(final String cql, final ResultSetExtractor rse, final Long timeout, - final TimeUnit timeUnit, final QueryOptions options) { - return rse.extractData(execute(new SessionCallback() { + public T queryAsynchronously(final String cql, final ResultSetExtractor resultSetExtractor, + final Long timeout, final TimeUnit timeUnit, final QueryOptions options) { + + return resultSetExtractor.extractData(execute(new SessionCallback() { @Override - public ResultSet doInSession(Session s) throws DataAccessException { - Statement statement = new SimpleStatement(cql); - addQueryOptions(statement, options); - ResultSetFuture rsf = s.executeAsync(statement); - ResultSet rs = null; + public ResultSet doInSession(Session session) throws DataAccessException { + Statement statement = addQueryOptions(new SimpleStatement(logCql(cql)), options); + + ResultSetFuture resultSetFuture = session.executeAsync(statement); + try { - rs = rsf.get(timeout, timeUnit); + return resultSetFuture.get(timeout, timeUnit); } catch (TimeoutException e) { - throw new QueryTimeoutException("Asyncronous Query Timed Out.", e); + throw new QueryTimeoutException(String.format( + "timeout occurred in [%1$d %2$s] while asynchronously executing CQL [%3$s]", + timeout, timeUnit, cql), e); } catch (InterruptedException e) { throw translateExceptionIfPossible(e); } catch (ExecutionException e) { @@ -266,61 +311,38 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } throw new CassandraUncategorizedDataAccessException("Unknown Throwable", e.getCause()); } - return rs; } })); } @Override - public ResultSetFuture queryAsynchronously(final String cql, final QueryOptions options) { + public ResultSetFuture queryAsynchronously(final String cql, final QueryOptions queryOptions) { return execute(new SessionCallback() { @Override - public ResultSetFuture doInSession(Session s) throws DataAccessException { - Statement statement = new SimpleStatement(cql); - addQueryOptions(statement, options); - return s.executeAsync(statement); + public ResultSetFuture doInSession(Session session) throws DataAccessException { + return session.executeAsync(addQueryOptions(new SimpleStatement(logCql(cql)), queryOptions)); } }); } @Override public Cancellable queryAsynchronously(String cql, Runnable listener) { - return queryAsynchronously(cql, listener, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); - } - }); + return queryAsynchronously(cql, listener, RUN_RUNNABLE_EXECUTOR); } @Override public Cancellable queryAsynchronously(String cql, AsynchronousQueryListener listener) { - return queryAsynchronously(cql, listener, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); - } - }); + return queryAsynchronously(cql, listener, RUN_RUNNABLE_EXECUTOR); } @Override - public Cancellable queryAsynchronously(String cql, Runnable listener, QueryOptions options) { - return queryAsynchronously(cql, listener, options, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); - } - }); + public Cancellable queryAsynchronously(String cql, Runnable listener, QueryOptions queryOptions) { + return queryAsynchronously(cql, listener, queryOptions, RUN_RUNNABLE_EXECUTOR); } @Override - public Cancellable queryAsynchronously(String cql, AsynchronousQueryListener listener, QueryOptions options) { - return queryAsynchronously(cql, listener, options, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); - } - }); + public Cancellable queryAsynchronously(String cql, AsynchronousQueryListener listener, QueryOptions queryOptions) { + return queryAsynchronously(cql, listener, queryOptions, RUN_RUNNABLE_EXECUTOR); } @Override @@ -334,82 +356,100 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } @Override - public Cancellable queryAsynchronously(final String cql, final Runnable listener, final QueryOptions options, + public Cancellable queryAsynchronously(final String cql, final Runnable listener, final QueryOptions queryOptions, final Executor executor) { + return execute(new SessionCallback() { @Override - public Cancellable doInSession(Session s) throws DataAccessException { - Statement statement = new SimpleStatement(cql); - addQueryOptions(statement, options); - ResultSetFuture rsf = s.executeAsync(statement); - rsf.addListener(listener, executor); - return new ResultSetFutureCancellable(rsf); + public Cancellable doInSession(Session session) throws DataAccessException { + Statement statement = addQueryOptions(new SimpleStatement(logCql("async execute CQL [{}]", cql)), + queryOptions); + + ResultSetFuture resultSetFuture = session.executeAsync(statement); + + resultSetFuture.addListener(listener, executor); + + return new ResultSetFutureCancellable(resultSetFuture); } }); } @Override public Cancellable queryAsynchronously(final String cql, final AsynchronousQueryListener listener, - final QueryOptions options, final Executor executor) { + final QueryOptions queryOptions, final Executor executor) { + return execute(new SessionCallback() { @Override - public Cancellable doInSession(Session s) throws DataAccessException { - Statement statement = new SimpleStatement(cql); - addQueryOptions(statement, options); - final ResultSetFuture rsf = s.executeAsync(statement); - Runnable wrapper = new Runnable() { + public Cancellable doInSession(Session session) throws DataAccessException { + Statement statement = addQueryOptions(new SimpleStatement(logCql("async execute CQL [{}]", cql)), + queryOptions); + + final ResultSetFuture resultSetFuture = session.executeAsync(statement); + + Runnable runnable = new Runnable() { @Override public void run() { - listener.onQueryComplete(rsf); + listener.onQueryComplete(resultSetFuture); } }; - rsf.addListener(wrapper, executor); - return new ResultSetFutureCancellable(rsf); + + resultSetFuture.addListener(runnable, executor); + + return new ResultSetFutureCancellable(resultSetFuture); } }); } - public T queryAsynchronously(final String cql, ResultSetFutureExtractor rse, final QueryOptions options) - throws DataAccessException { - return rse.extractData(execute(new SessionCallback() { + @SuppressWarnings("unused") + public T queryAsynchronously(String cql, ResultSetFutureExtractor resultSetFutureExtractor) + throws DataAccessException { + + return queryAsynchronously(cql, resultSetFutureExtractor, null); + } + + public T queryAsynchronously(final String cql, ResultSetFutureExtractor resultSetFutureExtractor, + final QueryOptions queryOptions) throws DataAccessException { + + return resultSetFutureExtractor.extractData(execute(new SessionCallback() { @Override - public ResultSetFuture doInSession(Session s) throws DataAccessException { - Statement statement = new SimpleStatement(cql); - addQueryOptions(statement, options); - return s.executeAsync(statement); + public ResultSetFuture doInSession(Session session) throws DataAccessException { + return session.executeAsync(addQueryOptions(new SimpleStatement( + logCql("async execute CQL [{}]", cql)), queryOptions)); } })); } - public T queryAsynchronously(final String cql, ResultSetFutureExtractor rse) throws DataAccessException { - return queryAsynchronously(cql, rse, null); + @Override + public T query(String cql, ResultSetExtractor resultSetExtractor) throws DataAccessException { + return query(cql, resultSetExtractor, null); } @Override - public T query(String cql, ResultSetExtractor rse) throws DataAccessException { - return query(cql, rse, null); + public T query(String cql, ResultSetExtractor resultSetExtractor, QueryOptions queryOptions) + throws DataAccessException { + + Assert.notNull(cql, "CQL must not be null"); + + return resultSetExtractor.extractData(doExecute(cql, queryOptions)); } @Override - public T query(String cql, ResultSetExtractor rse, QueryOptions options) throws DataAccessException { - Assert.notNull(cql); - ResultSet rs = doExecute(cql, options); - return rse.extractData(rs); + public void query(String cql, RowCallbackHandler rowCallbackHandler) throws DataAccessException { + query(cql, rowCallbackHandler, null); } @Override - public void query(String cql, RowCallbackHandler rch, QueryOptions options) throws DataAccessException { - process(doExecute(cql, options), rch); + public void query(String cql, RowCallbackHandler rowCallbackHandler, QueryOptions queryOptions) + throws DataAccessException { + + process(doExecute(cql, queryOptions), rowCallbackHandler); } @Override - public void query(String cql, RowCallbackHandler rch) throws DataAccessException { - query(cql, rch, null); - } + public List query(String cql, RowMapper rowMapper, QueryOptions queryOptions) + throws DataAccessException { - @Override - public List query(String cql, RowMapper rowMapper, QueryOptions options) throws DataAccessException { - return process(doExecute(cql, options), rowMapper); + return process(doExecute(cql, queryOptions), rowMapper); } @Override @@ -418,15 +458,8 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } @Override - public ResultSet query(String cql, QueryOptions options) { - - return query(cql, new ResultSetExtractor() { - - @Override - public ResultSet extractData(ResultSet rs) throws DriverException, DataAccessException { - return rs; - } - }, options); + public ResultSet query(String cql, QueryOptions queryOptions) { + return query(cql, RESULT_SET_RETURNING_EXTRACTOR, queryOptions); } @Override @@ -459,138 +492,89 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { return processOne(doExecute(cql, null), rowMapper); } - /** - * Execute a command at the Session Level - * - * @param callback - * @return - */ - protected T doExecute(SessionCallback callback) { - - Assert.notNull(callback); - - try { - - return callback.doInSession(getSession()); - - } catch (DataAccessException e) { - throw translateExceptionIfPossible(e); - } - } - + @SuppressWarnings("unused") protected ResultSet doExecute(String cql) { return doExecute(cql, null); } - protected ResultSet doExecute(String cql, QueryOptions options) { - return doExecute(addQueryOptions(new SimpleStatement(cql), options)); + protected ResultSet doExecute(String cql, QueryOptions queryOptions) { + return doExecute(addQueryOptions(new SimpleStatement(logCql(cql)), queryOptions)); } /** * Execute a command at the Session Level with optional options * - * @param q The query to execute. + * @param statement The query to execute. */ - protected ResultSet doExecute(final Statement q) { - + protected ResultSet doExecute(final Statement statement) { return doExecute(new SessionCallback() { - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - - if (log.isDebugEnabled()) { - log.debug("executing [{}]", q.toString()); - } - - return s.execute(q); + public ResultSet doInSession(Session session) throws DataAccessException { + logDebug("execute [{}]", statement); + return session.execute(statement); } }); } - protected ResultSetFuture doExecuteAsync(final Statement q) { - + protected ResultSetFuture doExecuteAsync(final Statement statement) { return doExecute(new SessionCallback() { - @Override - public ResultSetFuture doInSession(Session s) throws DataAccessException { - - if (log.isDebugEnabled()) { - log.debug("asynchronously executing [{}]", q.toString()); - } - return s.executeAsync(q); + public ResultSetFuture doInSession(Session session) throws DataAccessException { + logDebug("async execute [{}]", statement); + return session.executeAsync(statement); } }); } - protected Cancellable doExecuteAsync(final Statement q, final AsynchronousQueryListener listener) { - return doExecuteAsync(q, listener, null); + protected Cancellable doExecuteAsync(final Statement statement, final AsynchronousQueryListener listener) { + return doExecuteAsync(statement, listener, null); } - protected Cancellable doExecuteAsync(final Statement q, final AsynchronousQueryListener listener, - final QueryOptions options) { + protected Cancellable doExecuteAsync(final Statement statement, final AsynchronousQueryListener listener, + final QueryOptions queryOptions) { return doExecute(new SessionCallback() { - @Override - public Cancellable doInSession(Session s) throws DataAccessException { + public Cancellable doInSession(Session session) throws DataAccessException { + logDebug("async execute [{}]", statement); - if (log.isDebugEnabled()) { - log.debug("asynchronously executing [{}]", q.toString()); - } - - if (options != null) { - addQueryOptions(q, options); - } - - final ResultSetFuture rsf = s.executeAsync(q); + final ResultSetFuture resultSetFuture = session.executeAsync(addQueryOptions(statement, queryOptions)); if (listener != null) { - rsf.addListener(new Runnable() { - + resultSetFuture.addListener(new Runnable() { @Override public void run() { - listener.onQueryComplete(rsf); + listener.onQueryComplete(resultSetFuture); } - }, new Executor() { - - @Override - public void execute(Runnable command) { - command.run(); - } - }); + }, RUN_RUNNABLE_EXECUTOR); } - return new ResultSetFutureCancellable(rsf); + + return new ResultSetFutureCancellable(resultSetFuture); } }); } - /** - * @param row - * @return - */ protected Object firstColumnToObject(Row row) { - ColumnDefinitions cols = row.getColumnDefinitions(); - if (cols.size() == 0) { - return null; - } - return CodecRegistry.DEFAULT_INSTANCE.codecFor(cols.getType(0)).deserialize(row.getBytesUnsafe(0), ProtocolVersion.NEWEST_SUPPORTED); + Iterator columnDefinitions = row.getColumnDefinitions().iterator(); + return (columnDefinitions.hasNext() ? columnToObject(row, columnDefinitions.next()) : null); + } + + /* (non-Javadoc) */ + T columnToObject(Row row, Definition columnDefinition) { + TypeCodec typeCodec = CodecRegistry.DEFAULT_INSTANCE.codecFor(columnDefinition.getType()); + return typeCodec.deserialize(row.getBytesUnsafe(columnDefinition.getName()), ProtocolVersion.NEWEST_SUPPORTED); } - /** - * @param row - * @return - */ protected Map toMap(Row row) { - if (row == null) { - return null; - } + Map map = null; - ColumnDefinitions cols = row.getColumnDefinitions(); - Map map = new HashMap(cols.size()); + if (row != null) { + ColumnDefinitions columns = row.getColumnDefinitions(); + map = new HashMap(columns.size()); - for (Definition def : cols.asList()) { - String name = def.getName(); - map.put(name, CodecRegistry.DEFAULT_INSTANCE.codecFor(def.getType()).deserialize(row.getBytesUnsafe(name), ProtocolVersion.NEWEST_SUPPORTED)); + for (Definition columnDefinition : columns.asList()) { + map.put(columnDefinition.getName(), columnToObject(row, columnDefinition)); + } } return map; @@ -602,25 +586,20 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } /** - * Pulls the list of Hosts for the current Session - * - * @return + * Requests the set of hosts in the Cassandra cluster from the current {@link Session}. */ protected Set getHosts() { - return doExecute(new SessionCallback>() { - @Override - public Set doInSession(Session s) throws DataAccessException { - return s.getCluster().getMetadata().getAllHosts(); + public Set doInSession(Session session) throws DataAccessException { + return session.getCluster().getMetadata().getAllHosts(); } }); } @Override public Collection describeRing(HostMapper hostMapper) throws DataAccessException { - Set hosts = getHosts(); - return hostMapper.mapHosts(hosts); + return hostMapper.mapHosts(getHosts()); } @Override @@ -629,18 +608,13 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } @Override - public ResultSetFuture executeAsynchronously(final String cql, QueryOptions options) throws DataAccessException { - return doExecuteAsync(addQueryOptions(new SimpleStatement(cql), options)); + public ResultSetFuture executeAsynchronously(String cql, QueryOptions queryOptions) throws DataAccessException { + return doExecuteAsync(addQueryOptions(new SimpleStatement(logCql(cql)), queryOptions)); } @Override public Cancellable executeAsynchronously(String cql, Runnable listener) throws DataAccessException { - return executeAsynchronously(cql, listener, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); - } - }); + return executeAsynchronously(cql, listener, RUN_RUNNABLE_EXECUTOR); } @Override @@ -649,24 +623,20 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { return execute(new SessionCallback() { @Override - public Cancellable doInSession(Session s) throws DataAccessException { - Statement statement = new SimpleStatement(cql); - final ResultSetFuture rsf = s.executeAsync(statement); - rsf.addListener(listener, executor); - return new ResultSetFutureCancellable(rsf); + public Cancellable doInSession(Session session) throws DataAccessException { + Statement statement = new SimpleStatement(logCql("async execute CQL [{}]", cql)); + ResultSetFuture resultSetFuture = session.executeAsync(statement); + resultSetFuture.addListener(listener, executor); + return new ResultSetFutureCancellable(resultSetFuture); } }); } @Override - public Cancellable executeAsynchronously(String cql, AsynchronousQueryListener listener) throws DataAccessException { + public Cancellable executeAsynchronously(String cql, AsynchronousQueryListener listener) + throws DataAccessException { - return executeAsynchronously(cql, listener, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); - } - }); + return executeAsynchronously(cql, listener, RUN_RUNNABLE_EXECUTOR); } @Override @@ -675,141 +645,157 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { return execute(new SessionCallback() { @Override - public Cancellable doInSession(Session s) throws DataAccessException { - Statement statement = new SimpleStatement(cql); - final ResultSetFuture rsf = s.executeAsync(statement); - Runnable wrapper = new Runnable() { + public Cancellable doInSession(Session session) throws DataAccessException { + Statement statement = new SimpleStatement(logCql("async execute CQL [{}]", cql)); + + final ResultSetFuture resultSetFuture = session.executeAsync(statement); + + Runnable runnable = new Runnable() { @Override public void run() { - listener.onQueryComplete(rsf); + listener.onQueryComplete(resultSetFuture); } }; - rsf.addListener(wrapper, executor); - return new ResultSetFutureCancellable(rsf); + + resultSetFuture.addListener(runnable, executor); + + return new ResultSetFutureCancellable(resultSetFuture); } }); } @Override - public ResultSetFuture executeAsynchronously(Statement query) throws DataAccessException { - return doExecuteAsync(query); + public ResultSetFuture executeAsynchronously(Statement statement) throws DataAccessException { + return doExecuteAsync(statement); } @Override - public Cancellable executeAsynchronously(Statement query, Runnable listener) throws DataAccessException { - return executeAsynchronously(query, listener, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); - } - }); + public Cancellable executeAsynchronously(Statement statement, Runnable listener) throws DataAccessException { + return executeAsynchronously(statement, listener, RUN_RUNNABLE_EXECUTOR); } @Override - public Cancellable executeAsynchronously(Statement query, AsynchronousQueryListener listener) + public Cancellable executeAsynchronously(Statement statement, AsynchronousQueryListener listener) throws DataAccessException { - return executeAsynchronously(query, listener, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); - } - }); + + return executeAsynchronously(statement, listener, RUN_RUNNABLE_EXECUTOR); } @Override - public Cancellable executeAsynchronously(final Statement query, final Runnable listener, final Executor executor) - throws DataAccessException { - return execute(new SessionCallback() { - @Override - public Cancellable doInSession(Session s) throws DataAccessException { - final ResultSetFuture rsf = s.executeAsync(query); - rsf.addListener(listener, executor); - return new ResultSetFutureCancellable(rsf); - } - }); - } - - @Override - public Cancellable executeAsynchronously(final Statement query, final AsynchronousQueryListener listener, + public Cancellable executeAsynchronously(final Statement statement, final Runnable listener, final Executor executor) throws DataAccessException { return execute(new SessionCallback() { @Override - public Cancellable doInSession(Session s) throws DataAccessException { - final ResultSetFuture rsf = s.executeAsync(query); - if (listener != null) { - Runnable wrapper = new Runnable() { - @Override - public void run() { - listener.onQueryComplete(rsf); - } - }; - rsf.addListener(wrapper, executor); - } - return new ResultSetFutureCancellable(rsf); + public Cancellable doInSession(Session session) throws DataAccessException { + logDebug("executing [{}]", statement); + final ResultSetFuture resultSetFuture = session.executeAsync(statement); + resultSetFuture.addListener(listener, executor); + return new ResultSetFutureCancellable(resultSetFuture); } }); } @Override - public void process(ResultSet resultSet, RowCallbackHandler rch) throws DataAccessException { + public Cancellable executeAsynchronously(final Statement statement, final AsynchronousQueryListener listener, + final Executor executor) throws DataAccessException { + + return execute(new SessionCallback() { + @Override + public Cancellable doInSession(Session session) throws DataAccessException { + logDebug("executing [{}]", statement); + + final ResultSetFuture resultSetFuture = session.executeAsync(statement); + + Runnable runnable = new Runnable() { + @Override + public void run() { + listener.onQueryComplete(resultSetFuture); + } + }; + + resultSetFuture.addListener(runnable, executor); + + return new ResultSetFutureCancellable(resultSetFuture); + } + }); + } + + @Override + public void process(ResultSet resultSet, RowCallbackHandler rowCallbackHandler) throws DataAccessException { try { for (Row row : resultSet.all()) { - rch.processRow(row); + rowCallbackHandler.processRow(row); } - } catch (DriverException dx) { - throw translateExceptionIfPossible(dx); + } catch (DriverException e) { + throw translateExceptionIfPossible(e); } } @Override public List process(ResultSet resultSet, RowMapper rowMapper) throws DataAccessException { - List mappedRows = new ArrayList(); try { - int i = 0; - for (Row row : resultSet.all()) { - mappedRows.add(rowMapper.mapRow(row, i++)); + List rows = resultSet.all(); + List mappedRows = new ArrayList(rows.size()); + + int rowIndex = 0; + + for (Row row : rows) { + mappedRows.add(rowMapper.mapRow(row, rowIndex++)); } + + return mappedRows; } catch (DriverException dx) { throw translateExceptionIfPossible(dx); } - return mappedRows; } @Override public T processOne(ResultSet resultSet, RowMapper rowMapper) throws DataAccessException { - T row = null; Assert.notNull(resultSet, "ResultSet cannot be null"); + try { - List rows = resultSet.all(); - Assert.notNull(rows, "null row list returned from query"); - Assert.isTrue(rows.size() == 1, "row list has " + rows.size() + " rows instead of one"); - row = rowMapper.mapRow(rows.get(0), 0); - } catch (DriverException dx) { - throw translateExceptionIfPossible(dx); + Row row = resultSet.one(); + + if (row == null) { + throw new IncorrectResultSizeDataAccessException(1, 0); + } + + if (!resultSet.isExhausted()) { + throw new IncorrectResultSizeDataAccessException("ResultSet size exceeds 1", 1); + } + + return rowMapper.mapRow(row, 0); + } catch (DriverException e) { + throw translateExceptionIfPossible(e); } - return row; } - @SuppressWarnings("unchecked") @Override + @SuppressWarnings("unchecked") public T processOne(ResultSet resultSet, Class requiredType) throws DataAccessException { - if (resultSet == null) { - return null; + Assert.notNull(resultSet, "ResultSet cannot be null"); + + try { + Row row = resultSet.one(); + + if (row == null) { + throw new IncorrectResultSizeDataAccessException(1, 0); + } + + if (!resultSet.isExhausted()) { + throw new IncorrectResultSizeDataAccessException("ResultSet size exceeds 1", 1); + } + + return requiredType.cast(firstColumnToObject(row)); + } catch (DriverException e) { + throw translateExceptionIfPossible(e); } - Row row = resultSet.one(); - if (row == null) { - return null; - } - return (T) firstColumnToObject(row); } @Override public Map processMap(ResultSet resultSet) throws DataAccessException { - if (resultSet == null) { - return null; - } - return toMap(resultSet.one()); + return (resultSet != null ? toMap(resultSet.one()) : null); } @Override @@ -817,9 +803,11 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { public List processList(ResultSet resultSet, Class elementType) throws DataAccessException { List rows = resultSet.all(); List list = new ArrayList(rows.size()); + for (Row row : rows) { - list.add((T) firstColumnToObject(row)); + list.add(elementType.cast(firstColumnToObject(row))); } + return list; } @@ -827,123 +815,146 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { public List> processListOfMap(ResultSet resultSet) throws DataAccessException { List rows = resultSet.all(); List> list = new ArrayList>(rows.size()); + for (Row row : rows) { list.add(toMap(row)); } + return list; } /** - * Attempt to translate a Runtime Exception to a Spring Data Exception - * - * @param ex - * @return + * Attempts to translate the {@link RuntimeException} into a Spring Data {@link Exception}. */ - protected RuntimeException translateExceptionIfPossible(RuntimeException ex) { - RuntimeException resolved = getExceptionTranslator().translateExceptionIfPossible(ex); - return resolved == null ? ex : resolved; + @SuppressWarnings("all") + protected RuntimeException translateExceptionIfPossible(RuntimeException e) { + RuntimeException resolved = getExceptionTranslator().translateExceptionIfPossible(e); + return (resolved != null ? resolved : e); } - protected RuntimeException translateExceptionIfPossible(Exception ex) { - if (ex instanceof RuntimeException) { - return translateExceptionIfPossible((RuntimeException) ex); - } - return new CassandraUncategorizedDataAccessException("Caught Uncategorized Exception", ex); + @SuppressWarnings("all") + protected RuntimeException translateExceptionIfPossible(Exception e) { + return (e instanceof RuntimeException ? translateExceptionIfPossible((RuntimeException) e) : + new CassandraUncategorizedDataAccessException("Caught Uncategorized Exception", e)); } @Override - public T execute(PreparedStatementCreator psc, PreparedStatementCallback action) { + public T execute(PreparedStatementCreator preparedStatementCreator, + PreparedStatementCallback preparedStatementCallback) { try { - PreparedStatement ps = psc.createPreparedStatement(getSession()); - return action.doInPreparedStatement(ps); + PreparedStatement preparedStatement = preparedStatementCreator.createPreparedStatement(getSession()); + logDebug("executing [{}]", preparedStatement); + return preparedStatementCallback.doInPreparedStatement(preparedStatement); } catch (DriverException dx) { throw translateExceptionIfPossible(dx); } } @Override - public T execute(String cql, PreparedStatementCallback action) { - return execute(new CachedPreparedStatementCreator(cql), action); + public T execute(String cql, PreparedStatementCallback callback) { + return execute(new CachedPreparedStatementCreator(logCql(cql)), callback); } @Override - public T query(PreparedStatementCreator psc, ResultSetExtractor rse, QueryOptions options) - throws DataAccessException { - return query(psc, null, rse, options); + public T query(PreparedStatementCreator preparedStatementCreator, ResultSetExtractor resultSetExtractor) + throws DataAccessException { + + return query(preparedStatementCreator, resultSetExtractor, null); } @Override - public T query(PreparedStatementCreator psc, ResultSetExtractor rse) throws DataAccessException { - return query(psc, rse, null); + public T query(PreparedStatementCreator preparedStatementCreator, ResultSetExtractor resultSetExtractor, + QueryOptions queryOptions) throws DataAccessException { + + return query(preparedStatementCreator, null, resultSetExtractor, queryOptions); } @Override - public void query(PreparedStatementCreator psc, RowCallbackHandler rch, QueryOptions options) - throws DataAccessException { - query(psc, null, rch, options); + public void query(PreparedStatementCreator preparedStatementCreator, RowCallbackHandler rowCallbackHandler) + throws DataAccessException { + + query(preparedStatementCreator, rowCallbackHandler, null); } @Override - public void query(PreparedStatementCreator psc, RowCallbackHandler rch) throws DataAccessException { - query(psc, rch, null); + public void query(PreparedStatementCreator preparedStatementCreator, RowCallbackHandler rowCallbackHandler, + QueryOptions queryOptions) throws DataAccessException { + + query(preparedStatementCreator, null, rowCallbackHandler, queryOptions); } @Override - public List query(PreparedStatementCreator psc, RowMapper rowMapper, QueryOptions options) - throws DataAccessException { - return query(psc, null, rowMapper, options); + public List query(PreparedStatementCreator preparedStatementCreator, RowMapper rowMapper) + throws DataAccessException { + + return query(preparedStatementCreator, rowMapper, null); } @Override - public List query(PreparedStatementCreator psc, RowMapper rowMapper) throws DataAccessException { - return query(psc, rowMapper, null); + public List query(PreparedStatementCreator preparedStatementCreator, RowMapper rowMapper, + QueryOptions queryOptions) throws DataAccessException { + + return query(preparedStatementCreator, null, rowMapper, queryOptions); } @Override - public T query(String cql, PreparedStatementBinder psb, ResultSetExtractor rse, QueryOptions options) - throws DataAccessException { - return query(new CachedPreparedStatementCreator(cql), psb, rse, options); + public T query(String cql, PreparedStatementBinder preparedStatementBinder, + ResultSetExtractor resultSetExtractor) throws DataAccessException { + + return query(cql, preparedStatementBinder, resultSetExtractor, null); } @Override - public T query(String cql, PreparedStatementBinder psb, ResultSetExtractor rse) throws DataAccessException { - return query(cql, psb, rse, null); + public T query(String cql, PreparedStatementBinder preparedStatementBinder, + ResultSetExtractor resultSetExtractor, QueryOptions queryOptions) throws DataAccessException { + + return query(new CachedPreparedStatementCreator(logCql(cql)), + preparedStatementBinder, resultSetExtractor, queryOptions); } @Override - public void query(String cql, PreparedStatementBinder psb, RowCallbackHandler rch, QueryOptions options) - throws DataAccessException { - query(new CachedPreparedStatementCreator(cql), psb, rch, options); + public void query(String cql, PreparedStatementBinder preparedStatementBinder, + RowCallbackHandler rowCallbackHandler) throws DataAccessException { + + query(cql, preparedStatementBinder, rowCallbackHandler, null); } @Override - public void query(String cql, PreparedStatementBinder psb, RowCallbackHandler rch) throws DataAccessException { - query(cql, psb, rch, null); + public void query(String cql, PreparedStatementBinder preparedStatementBinder, + RowCallbackHandler rowCallbackHandler, QueryOptions queryOptions) throws DataAccessException { + + query(new CachedPreparedStatementCreator(logCql(cql)), preparedStatementBinder, rowCallbackHandler, + queryOptions); } @Override - public List query(String cql, PreparedStatementBinder psb, RowMapper rowMapper, QueryOptions options) - throws DataAccessException { - return query(new CachedPreparedStatementCreator(cql), psb, rowMapper, options); + public List query(String cql, PreparedStatementBinder preparedStatementBinder, RowMapper rowMapper) + throws DataAccessException { + + return query(cql, preparedStatementBinder, rowMapper, null); } @Override - public List query(String cql, PreparedStatementBinder psb, RowMapper rowMapper) throws DataAccessException { - return query(cql, psb, rowMapper, null); + public List query(String cql, PreparedStatementBinder preparedStatementBinder, RowMapper rowMapper, + QueryOptions queryOptions) throws DataAccessException { + + return query(new CachedPreparedStatementCreator(logCql(cql)), preparedStatementBinder, rowMapper, queryOptions); } @Override public void ingest(String cql, RowIterator rowIterator, WriteOptions options) { - CachedPreparedStatementCreator cpsc = new CachedPreparedStatementCreator(cql); + CachedPreparedStatementCreator cachedPreparedStatementCreator = + new CachedPreparedStatementCreator(logCql(cql)); - PreparedStatement preparedStatement = cpsc.createPreparedStatement(getSession()); - addPreparedStatementOptions(preparedStatement, options); + PreparedStatement preparedStatement = addPreparedStatementOptions( + cachedPreparedStatementCreator.createPreparedStatement(getSession()), options); + + Session session = getSession(); - Session s = getSession(); while (rowIterator.hasNext()) { - s.executeAsync(preparedStatement.bind(rowIterator.next())); + session.executeAsync(preparedStatement.bind(rowIterator.next())); } } @@ -952,58 +963,57 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { ingest(cql, rowIterator, null); } - @Override - public void ingest(String cql, final List> rows, WriteOptions options) { - - Assert.notNull(rows); - Assert.notEmpty(rows); - - ingest(cql, new RowIterator() { - - Iterator> i = rows.iterator(); - - @Override - public Object[] next() { - return i.next().toArray(); - } - - @Override - public boolean hasNext() { - return i.hasNext(); - } - - }, options); - - } - @Override public void ingest(String cql, List> rows) { ingest(cql, rows, null); } @Override - public void ingest(String cql, final Object[][] rows, WriteOptions options) { + public void ingest(String cql, final List> rows, WriteOptions writeOptions) { + Assert.notNull(rows); + Assert.notEmpty(rows); + ingest(cql, new RowIterator() { + + Iterator> rowIterator = rows.iterator(); + + @Override + public Object[] next() { + return rowIterator.next().toArray(); + } + + @Override + public boolean hasNext() { + return rowIterator.hasNext(); + } + + }, writeOptions); + } + + @Override + public void ingest(String cql, Object[][] rows) { + ingest(cql, rows, null); + } + + @Override + public void ingest(String cql, final Object[][] rows, WriteOptions writeOptions) { ingest(cql, new RowIterator() { int index = 0; @Override public Object[] next() { + if (!hasNext()) { + throw new NoSuchElementException("No more elements"); + } return rows[index++]; } @Override public boolean hasNext() { - return index < rows.length; + return (index < rows.length); } - - }, options); - } - - @Override - public void ingest(String cql, final Object[][] rows) { - ingest(cql, rows, null); + }, writeOptions); } @Override @@ -1013,194 +1023,176 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public void truncate(CqlIdentifier tableName) throws DataAccessException { - Truncate truncate = QueryBuilder.truncate(tableName.toCql()); - doExecute(truncate); + doExecute(QueryBuilder.truncate(logCql(tableName.toCql()))); } @Override - public T query(PreparedStatementCreator psc, final PreparedStatementBinder psb, final ResultSetExtractor rse, - final QueryOptions options) throws DataAccessException { + public T query(PreparedStatementCreator preparedStatementCreator, + PreparedStatementBinder preparedStatementBinder, ResultSetExtractor resultSetExtractor) + throws DataAccessException { - Assert.notNull(rse, "ResultSetExtractor must not be null"); - logger.debug("Executing prepared CQL query"); + return query(preparedStatementCreator, preparedStatementBinder, resultSetExtractor, null); + } - return execute(psc, new PreparedStatementCallback() { + @Override + public T query(PreparedStatementCreator preparedStatementCreator, + final PreparedStatementBinder preparedStatementBinder, final ResultSetExtractor resultSetExtractor, + final QueryOptions queryOptions) throws DataAccessException { + + Assert.notNull(resultSetExtractor, "ResultSetExtractor must not be null"); + + return execute(preparedStatementCreator, new PreparedStatementCallback() { @Override - public T doInPreparedStatement(PreparedStatement ps) throws DriverException { - ResultSet rs = null; - BoundStatement bs = null; - if (psb != null) { - bs = psb.bindValues(ps); - } else { - bs = ps.bind(); - } - rs = doExecute(addQueryOptions(bs, options)); - return rse.extractData(rs); + public T doInPreparedStatement(PreparedStatement preparedStatement) throws DriverException { + BoundStatement boundStatement = (preparedStatementBinder != null + ? preparedStatementBinder.bindValues(preparedStatement) + : preparedStatement.bind()); + + return resultSetExtractor.extractData(doExecute(addQueryOptions(boundStatement, queryOptions))); } }); } @Override - public T query(PreparedStatementCreator psc, PreparedStatementBinder psb, ResultSetExtractor rse) - throws DataAccessException { - return query(psc, psb, rse, null); - } + public void query(PreparedStatementCreator preparedStatementCreator, + final PreparedStatementBinder preparedStatementBinder, final RowCallbackHandler rowCallbackHandler, + final QueryOptions queryOptions) throws DataAccessException { - @Override - public void query(PreparedStatementCreator psc, final PreparedStatementBinder psb, final RowCallbackHandler rch, - final QueryOptions options) throws DataAccessException { + Assert.notNull(rowCallbackHandler, "RowCallbackHandler must not be null"); - Assert.notNull(rch, "RowCallbackHandler must not be null"); - logger.debug("Executing prepared CQL query"); - - execute(psc, new PreparedStatementCallback() { + execute(preparedStatementCreator, new PreparedStatementCallback() { @Override - public Object doInPreparedStatement(PreparedStatement ps) throws DriverException { - ResultSet rs = null; - BoundStatement bs = null; - if (psb != null) { - bs = psb.bindValues(ps); - } else { - bs = ps.bind(); - } - rs = doExecute(addQueryOptions(bs, options)); - process(rs, rch); + public Object doInPreparedStatement(PreparedStatement preparedStatement) throws DriverException { + BoundStatement boundStatement = (preparedStatementBinder != null + ? preparedStatementBinder.bindValues(preparedStatement) + : preparedStatement.bind()); + + process(doExecute(addQueryOptions(boundStatement, queryOptions)), rowCallbackHandler); + return null; } }); } @Override - public void query(PreparedStatementCreator psc, PreparedStatementBinder psb, RowCallbackHandler rch) - throws DataAccessException { - query(psc, psb, rch, null); + public void query(PreparedStatementCreator preparedStatementCreator, + PreparedStatementBinder preparedStatementBinder, RowCallbackHandler rowCallbackHandler) + throws DataAccessException { + + query(preparedStatementCreator, preparedStatementBinder, rowCallbackHandler, null); } @Override - public List query(PreparedStatementCreator psc, final PreparedStatementBinder psb, - final RowMapper rowMapper, final QueryOptions options) throws DataAccessException { + public List query(PreparedStatementCreator preparedStatementCreator, + final PreparedStatementBinder preparedStatementBinder, final RowMapper rowMapper, + final QueryOptions queryOptions) throws DataAccessException { + Assert.notNull(rowMapper, "RowMapper must not be null"); - logger.debug("Executing prepared CQL query"); - return execute(psc, new PreparedStatementCallback>() { + return execute(preparedStatementCreator, new PreparedStatementCallback>() { @Override - public List doInPreparedStatement(PreparedStatement ps) throws DriverException { - ResultSet rs = null; - BoundStatement bs = null; - if (psb != null) { - bs = psb.bindValues(ps); - } else { - bs = ps.bind(); - } - rs = doExecute(addQueryOptions(bs, options)); + public List doInPreparedStatement(PreparedStatement preparedStatement) throws DriverException { + BoundStatement boundStatement = (preparedStatementBinder != null + ? preparedStatementBinder.bindValues(preparedStatement) + : preparedStatement.bind()); - return process(rs, rowMapper); + return process(doExecute(addQueryOptions(boundStatement, queryOptions)), rowMapper); } }); } @Override - public List query(PreparedStatementCreator psc, PreparedStatementBinder psb, RowMapper rowMapper) - throws DataAccessException { - return query(psc, psb, rowMapper, null); + public List query(PreparedStatementCreator preparedStatementCreator, + PreparedStatementBinder preparedStatementBinder, RowMapper rowMapper) throws DataAccessException { + + return query(preparedStatementCreator, preparedStatementBinder, rowMapper, null); } @Override - public ResultSet execute(final DropTableSpecification specification) { - + public ResultSet execute(final AlterKeyspaceSpecification specification) { return execute(new SessionCallback() { - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(DropTableCqlGenerator.toCql(specification)); - } - }); - } - - @Override - public ResultSet execute(final CreateTableSpecification specification) { - - return execute(new SessionCallback() { - - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(CreateTableCqlGenerator.toCql(specification)); - } - }); - } - - @Override - public ResultSet execute(final AlterTableSpecification specification) { - - return execute(new SessionCallback() { - - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(AlterTableCqlGenerator.toCql(specification)); - } - }); - } - - @Override - public ResultSet execute(final DropKeyspaceSpecification specification) { - - return execute(new SessionCallback() { - - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(DropKeyspaceCqlGenerator.toCql(specification)); + public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logCql(AlterKeyspaceCqlGenerator.toCql(specification))); } }); } @Override public ResultSet execute(final CreateKeyspaceSpecification specification) { - return execute(new SessionCallback() { - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(CreateKeyspaceCqlGenerator.toCql(specification)); + public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logCql(CreateKeyspaceCqlGenerator.toCql(specification))); } }); } @Override - public ResultSet execute(final AlterKeyspaceSpecification specification) { - + public ResultSet execute(final DropKeyspaceSpecification specification) { return execute(new SessionCallback() { - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(AlterKeyspaceCqlGenerator.toCql(specification)); + public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logCql(DropKeyspaceCqlGenerator.toCql(specification))); } }); } @Override - public ResultSet execute(final DropIndexSpecification specification) { - + public ResultSet execute(final AlterTableSpecification specification) { return execute(new SessionCallback() { - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(DropIndexCqlGenerator.toCql(specification)); + public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logCql(AlterTableCqlGenerator.toCql(specification))); + } + }); + } + + @Override + public ResultSet execute(final CreateTableSpecification specification) { + return execute(new SessionCallback() { + @Override + public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logCql(CreateTableCqlGenerator.toCql(specification))); + } + }); + } + + @Override + public ResultSet execute(final DropTableSpecification specification) { + return execute(new SessionCallback() { + @Override + public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logCql(DropTableCqlGenerator.toCql(specification))); } }); } @Override public ResultSet execute(final CreateIndexSpecification specification) { - return execute(new SessionCallback() { - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(CreateIndexCqlGenerator.toCql(specification)); + public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logCql(CreateIndexCqlGenerator.toCql(specification))); } }); } + @Override + public ResultSet execute(final DropIndexSpecification specification) { + return execute(new SessionCallback() { + @Override + public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logCql(DropIndexCqlGenerator.toCql(specification))); + } + }); + } + + @Override + public void execute(Batch batch) throws DataAccessException { + doExecute(batch); + } + @Override public void execute(Delete delete) throws DataAccessException { doExecute(delete); @@ -1211,24 +1203,14 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { doExecute(insert); } - @Override - public void execute(Update update) throws DataAccessException { - doExecute(update); - } - - @Override - public void execute(Batch batch) throws DataAccessException { - doExecute(batch); - } - @Override public void execute(Truncate truncate) throws DataAccessException { doExecute(truncate); } @Override - public long count(CqlIdentifier tableName) { - return selectCount(QueryBuilder.select().countAll().from(tableName.toCql())); + public void execute(Update update) throws DataAccessException { + doExecute(update); } @Override @@ -1236,16 +1218,20 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { return count(cqlId(tableName)); } - protected long selectCount(Select select) { + @Override + public long count(CqlIdentifier tableName) { + return selectCount(QueryBuilder.select().countAll().from(tableName.toCql())); + } + protected long selectCount(final Select select) { return query(select, new ResultSetExtractor() { - @Override - public Long extractData(ResultSet rs) throws DriverException, DataAccessException { + public Long extractData(ResultSet resultSet) throws DriverException, DataAccessException { + Row row = resultSet.one(); - Row row = rs.one(); if (row == null) { - throw new InvalidDataAccessApiUsageException(String.format("count query did not return any results")); + throw new InvalidDataAccessApiUsageException(String.format( + "count query [%1$s] did not return any results", select)); } return row.getLong(0); @@ -1254,8 +1240,8 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } @Override - public ResultSetFuture executeAsynchronously(Truncate truncate) throws DataAccessException { - return doExecuteAsync(truncate); + public ResultSetFuture executeAsynchronously(Batch batch) throws DataAccessException { + return doExecuteAsync(batch); } @Override @@ -1268,73 +1254,65 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { return doExecuteAsync(insert); } + @Override + public ResultSetFuture executeAsynchronously(Truncate truncate) throws DataAccessException { + return doExecuteAsync(truncate); + } + @Override public ResultSetFuture executeAsynchronously(Update update) throws DataAccessException { return doExecuteAsync(update); } @Override - public ResultSetFuture executeAsynchronously(Batch batch) throws DataAccessException { - return doExecuteAsync(batch); - } - - @Override - public Cancellable executeAsynchronously(Truncate truncate, AsynchronousQueryListener listener) + public Cancellable executeAsynchronously(Batch batch, AsynchronousQueryListener listener) throws DataAccessException { - return doExecuteAsync(truncate, listener); + + return doExecuteAsync(batch, listener); } @Override public Cancellable executeAsynchronously(Delete delete, AsynchronousQueryListener listener) throws DataAccessException { + return doExecuteAsync(delete, listener); } @Override public Cancellable executeAsynchronously(Insert insert, AsynchronousQueryListener listener) throws DataAccessException { + return doExecuteAsync(insert, listener); } + @Override + public Cancellable executeAsynchronously(Truncate truncate, AsynchronousQueryListener listener) + throws DataAccessException { + + return doExecuteAsync(truncate, listener); + } + @Override public Cancellable executeAsynchronously(Update update, AsynchronousQueryListener listener) throws DataAccessException { - return doExecuteAsync(update, listener); - } - @Override - public Cancellable executeAsynchronously(Batch batch, AsynchronousQueryListener listener) throws DataAccessException { - return doExecuteAsync(batch, listener); + return doExecuteAsync(update, listener); } @Override public ResultSetFuture queryAsynchronously(final Select select) { return execute(new SessionCallback() { @Override - public ResultSetFuture doInSession(Session s) throws DataAccessException { - return s.executeAsync(select); - } - }); - } - - @Override - public Cancellable queryAsynchronously(Select select, Runnable listener) { - return queryAsynchronously(select, listener, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); + public ResultSetFuture doInSession(Session session) throws DataAccessException { + logDebug("async query [{}]", select); + return session.executeAsync(select); } }); } @Override public Cancellable queryAsynchronously(Select select, AsynchronousQueryListener listener) { - return queryAsynchronously(select, listener, new Executor() { - @Override - public void execute(Runnable command) { - command.run(); - } - }); + return queryAsynchronously(select, listener, RUN_RUNNABLE_EXECUTOR); } @Override @@ -1343,52 +1321,57 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { return execute(new SessionCallback() { @Override - public Cancellable doInSession(Session s) throws DataAccessException { - final ResultSetFuture rsf = s.executeAsync(select); + public Cancellable doInSession(Session session) throws DataAccessException { + logDebug("async query [{}]", select); + + final ResultSetFuture resultSetFuture = session.executeAsync(select); + Runnable wrapper = new Runnable() { @Override public void run() { - listener.onQueryComplete(rsf); + listener.onQueryComplete(resultSetFuture); } }; - rsf.addListener(wrapper, executor); - return new ResultSetFutureCancellable(rsf); + + resultSetFuture.addListener(wrapper, executor); + + return new ResultSetFutureCancellable(resultSetFuture); } }); } + @Override + public Cancellable queryAsynchronously(Select select, Runnable listener) { + return queryAsynchronously(select, listener, RUN_RUNNABLE_EXECUTOR); + } + @Override public Cancellable queryAsynchronously(final Select select, final Runnable listener, final Executor executor) { return execute(new SessionCallback() { @Override - public Cancellable doInSession(Session s) throws DataAccessException { - ResultSetFuture rsf = s.executeAsync(select); - rsf.addListener(listener, executor); - return new ResultSetFutureCancellable(rsf); + public Cancellable doInSession(Session session) throws DataAccessException { + logDebug("async query [{}]", select); + ResultSetFuture resultSetFuture = session.executeAsync(select); + resultSetFuture.addListener(listener, executor); + return new ResultSetFutureCancellable(resultSetFuture); } }); } @Override public ResultSet query(Select select) { - return query(select, new ResultSetExtractor() { - @Override - public ResultSet extractData(ResultSet rs) throws DriverException, DataAccessException { - return rs; - } - }); + return query(select, RESULT_SET_RETURNING_EXTRACTOR); } @Override - public T query(Select select, ResultSetExtractor rse) throws DataAccessException { + public T query(Select select, ResultSetExtractor resultSetExtractor) throws DataAccessException { Assert.notNull(select); - ResultSet rs = doExecute(select); - return rse.extractData(rs); + return resultSetExtractor.extractData(doExecute(select)); } @Override - public void query(Select select, RowCallbackHandler rch) throws DataAccessException { - process(doExecute(select), rch); + public void query(Select select, RowCallbackHandler rowCallbackHandler) throws DataAccessException { + process(doExecute(select), rowCallbackHandler); } @Override @@ -1422,19 +1405,18 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } @Override - public Cancellable queryForListAsynchronously(Select select, final Class elementType, + public Cancellable queryForListAsynchronously(Select select, final Class requiredType, final QueryForListListener listener) throws DataAccessException { - Assert.notNull(select); - Assert.notNull(elementType); - Assert.notNull(listener); + Assert.notNull(select, "Select cannot be null"); + Assert.notNull(requiredType, "Required type cannot be null"); + Assert.notNull(listener, "Listener cannot be null"); return doExecuteAsync(select, new AsynchronousQueryListener() { - @Override - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - listener.onQueryComplete(processList(rsf.getUninterruptibly(), elementType)); + listener.onQueryComplete(processList(resultSetFuture.getUninterruptibly(), requiredType)); } catch (Exception e) { listener.onException(translateExceptionIfPossible(e)); } @@ -1443,19 +1425,18 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } @Override - public Cancellable queryForListAsynchronously(String select, final Class elementType, + public Cancellable queryForListAsynchronously(String select, final Class requiredType, final QueryForListListener listener) throws DataAccessException { - Assert.hasText(select); - Assert.notNull(elementType); - Assert.notNull(listener); - - return doExecuteAsync(new SimpleStatement(select), new AsynchronousQueryListener() { + Assert.hasText(select, "Select cannot be null"); + Assert.notNull(requiredType, "Required type cannot be null"); + Assert.notNull(listener, "Listener cannot be null"); + return doExecuteAsync(new SimpleStatement(logCql(select)), new AsynchronousQueryListener() { @Override - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - listener.onQueryComplete(processList(rsf.getUninterruptibly(), elementType)); + listener.onQueryComplete(processList(resultSetFuture.getUninterruptibly(), requiredType)); } catch (Exception e) { listener.onException(translateExceptionIfPossible(e)); } @@ -1468,11 +1449,10 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { final QueryForListListener> listener) throws DataAccessException { return doExecuteAsync(select, new AsynchronousQueryListener() { - @Override - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - listener.onQueryComplete(processListOfMap(rsf.getUninterruptibly())); + listener.onQueryComplete(processListOfMap(resultSetFuture.getUninterruptibly())); } catch (Exception e) { listener.onException(translateExceptionIfPossible(e)); } @@ -1483,15 +1463,16 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public Cancellable queryForListOfMapAsynchronously(String cql, final QueryForListListener> listener) throws DataAccessException { + return queryForListOfMapAsynchronously(cql, listener, null); } @Override public Cancellable queryForListOfMapAsynchronously(String cql, - final QueryForListListener> listener, QueryOptions options) throws DataAccessException { - - return doExecuteAsync(new SimpleStatement(cql), new AsynchronousQueryListener() { + final QueryForListListener> listener, QueryOptions queryOptions) + throws DataAccessException { + return doExecuteAsync(new SimpleStatement(logCql(cql)), new AsynchronousQueryListener() { @Override public void onQueryComplete(ResultSetFuture rsf) { try { @@ -1500,29 +1481,30 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { listener.onException(translateExceptionIfPossible(e)); } } - }, options); + }, queryOptions); } @Override - public Cancellable queryForMapAsynchronously(String cql, QueryForMapListener listener) throws DataAccessException { + public Cancellable queryForMapAsynchronously(String cql, QueryForMapListener listener) + throws DataAccessException { + return queryForMapAsynchronously(cql, listener, null); } @Override public Cancellable queryForMapAsynchronously(String cql, final QueryForMapListener listener, - final QueryOptions options) throws DataAccessException { - - return doExecuteAsync(new SimpleStatement(cql), new AsynchronousQueryListener() { + final QueryOptions queryOptions) throws DataAccessException { + return doExecuteAsync(new SimpleStatement(logCql(cql)), new AsynchronousQueryListener() { @Override - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - listener.onQueryComplete(processMap(rsf.getUninterruptibly())); + listener.onQueryComplete(processMap(resultSetFuture.getUninterruptibly())); } catch (Exception e) { listener.onException(translateExceptionIfPossible(e)); } } - }, options); + }, queryOptions); } @Override @@ -1530,11 +1512,10 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { throws DataAccessException { return doExecuteAsync(select, new AsynchronousQueryListener() { - @Override - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - listener.onQueryComplete(processMap(rsf.getUninterruptibly())); + listener.onQueryComplete(processMap(resultSetFuture.getUninterruptibly())); } catch (Exception e) { listener.onException(translateExceptionIfPossible(e)); } @@ -1547,11 +1528,10 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { final QueryForObjectListener listener) throws DataAccessException { return doExecuteAsync(select, new AsynchronousQueryListener() { - @Override - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - listener.onQueryComplete(processOne(rsf.getUninterruptibly(), requiredType)); + listener.onQueryComplete(processOne(resultSetFuture.getUninterruptibly(), requiredType)); } catch (Exception e) { listener.onException(translateExceptionIfPossible(e)); } @@ -1562,6 +1542,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public Cancellable queryForObjectAsynchronously(String cql, Class requiredType, QueryForObjectListener listener) throws DataAccessException { + return queryForObjectAsynchronously(cql, requiredType, listener, null); } @@ -1569,12 +1550,11 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { public Cancellable queryForObjectAsynchronously(String cql, final Class requiredType, final QueryForObjectListener listener, QueryOptions options) throws DataAccessException { - return doExecuteAsync(new SimpleStatement(cql), new AsynchronousQueryListener() { - + return doExecuteAsync(new SimpleStatement(logCql(cql)), new AsynchronousQueryListener() { @Override - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - listener.onQueryComplete(processOne(rsf.getUninterruptibly(), requiredType)); + listener.onQueryComplete(processOne(resultSetFuture.getUninterruptibly(), requiredType)); } catch (Exception e) { listener.onException(translateExceptionIfPossible(e)); } @@ -1585,6 +1565,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public Cancellable queryForObjectAsynchronously(String cql, RowMapper rowMapper, QueryForObjectListener listener) throws DataAccessException { + return queryForObjectAsynchronously(cql, rowMapper, listener, null); } @@ -1592,12 +1573,11 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { public Cancellable queryForObjectAsynchronously(String cql, final RowMapper rowMapper, final QueryForObjectListener listener, QueryOptions options) throws DataAccessException { - return doExecuteAsync(new SimpleStatement(cql), new AsynchronousQueryListener() { - + return doExecuteAsync(new SimpleStatement(logCql(cql)), new AsynchronousQueryListener() { @Override - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - listener.onQueryComplete(processOne(rsf.getUninterruptibly(), rowMapper)); + listener.onQueryComplete(processOne(resultSetFuture.getUninterruptibly(), rowMapper)); } catch (Exception e) { listener.onException(translateExceptionIfPossible(e)); } @@ -1610,11 +1590,10 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { final QueryForObjectListener listener) throws DataAccessException { return doExecuteAsync(select, new AsynchronousQueryListener() { - @Override - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - listener.onQueryComplete(processOne(rsf.getUninterruptibly(), rowMapper)); + listener.onQueryComplete(processOne(resultSetFuture.getUninterruptibly(), rowMapper)); } catch (Exception e) { listener.onException(translateExceptionIfPossible(e)); } @@ -1623,22 +1602,24 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { } @Override - public ResultSet getResultSetUninterruptibly(ResultSetFuture rsf) { - return getResultSetUninterruptibly(rsf, 0, null); + public ResultSet getResultSetUninterruptibly(ResultSetFuture resultSetFuture) { + return getResultSetUninterruptibly(resultSetFuture, 0, null); } @Override - public ResultSet getResultSetUninterruptibly(ResultSetFuture rsf, long millis) { - return getResultSetUninterruptibly(rsf, millis, TimeUnit.MILLISECONDS); + public ResultSet getResultSetUninterruptibly(ResultSetFuture resultSetFuture, long milliseconds) { + return getResultSetUninterruptibly(resultSetFuture, milliseconds, TimeUnit.MILLISECONDS); } @Override - public ResultSet getResultSetUninterruptibly(ResultSetFuture rsf, long timeout, TimeUnit unit) { + public ResultSet getResultSetUninterruptibly(ResultSetFuture resultSetFuture, long timeout, TimeUnit timeUnit) { try { - return timeout <= 0 ? rsf.getUninterruptibly() : rsf.getUninterruptibly(timeout, - unit == null ? TimeUnit.MILLISECONDS : unit); - } catch (Exception x) { - throw translateExceptionIfPossible(x); + timeUnit = (timeUnit != null ? timeUnit : TimeUnit.MILLISECONDS); + + return (timeout > 0 ? resultSetFuture.getUninterruptibly(timeout, timeUnit) + : resultSetFuture.getUninterruptibly()); + } catch (Exception e) { + throw translateExceptionIfPossible(e); } } } diff --git a/spring-cql/src/main/java/org/springframework/cassandra/support/CassandraAccessor.java b/spring-cql/src/main/java/org/springframework/cassandra/support/CassandraAccessor.java index d96f1f8de..f0eebdab3 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/support/CassandraAccessor.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/support/CassandraAccessor.java @@ -23,70 +23,84 @@ import org.springframework.util.Assert; import com.datastax.driver.core.Session; /** - * A {@link CassandraAccessor} is able to access a Cassandra {@link Session} and the - * {@link CassandraExceptionTranslator}. + * {@link CassandraAccessor} provides access to a Cassandra {@link Session} and the {@link CassandraExceptionTranslator} + * . *

* Classes providing a higher abstraction level usually extend {@link CassandraAccessor} to provide a richer set of * functionality on top of a {@link Session}. * * @author David Webb * @author Mark Paluch + * @author John Blum * @see org.springframework.beans.factory.InitializingBean + * @see com.datastax.driver.core.Session */ public class CassandraAccessor implements InitializingBean { - /** Logger available to subclasses */ + CassandraExceptionTranslator exceptionTranslator = new CassandraExceptionTranslator(); + protected final Logger logger = LoggerFactory.getLogger(getClass()); private Session session; - private CassandraExceptionTranslator exceptionTranslator = new CassandraExceptionTranslator(); /** - * Ensure that the Cassandra Session has been set + * Ensures the Cassandra {@link Session} and exception translator has been propertly set. */ @Override public void afterPropertiesSet() { + Assert.state(session != null, "Session must not be null"); + } - Assert.notNull(session, "Session must not be null!"); - Assert.notNull(exceptionTranslator, "CassandraExceptionTranslator must not be null!"); + /* (non-Javadoc) */ + protected void logDebug(String logMessage, Object... array) { + if (logger.isDebugEnabled()) { + logger.debug(logMessage, array); + } } /** - * Set the exception translator for this instance. + * Sets the exception translator used by this template to translate Cassandra specific Exceptions into Spring DAO's + * Exception Hierarchy. * - * @param exceptionTranslator the exception translator to set, must not be {@literal null}. + * @param exceptionTranslator exception translator to set; must not be {@literal null}. * @see org.springframework.cassandra.support.CassandraExceptionTranslator */ public void setExceptionTranslator(CassandraExceptionTranslator exceptionTranslator) { - - Assert.notNull(exceptionTranslator, "CassandraExceptionTranslator must not be null!"); + Assert.notNull(exceptionTranslator, "CassandraExceptionTranslator must not be null"); this.exceptionTranslator = exceptionTranslator; } /** - * Return the exception translator for this instance. + * Return the exception translator used by this template to translate Cassandra specific Exceptions into Spring DAO's + * Exception Hierarchy. * - * @return the exception translator + * @return the Cassandra exception translator. + * @see org.springframework.cassandra.support.CassandraExceptionTranslator */ public CassandraExceptionTranslator getExceptionTranslator() { + Assert.state(this.exceptionTranslator != null, "CassandraExceptionTranslator was not properly initialized"); return this.exceptionTranslator; } /** - * Returns the session. + * Sets the Cassandra {@link Session} used by this template to perform Cassandra data access operations. * - * @return the session. + * @param session Cassandra {@link Session} used by this template. Must not be{@literal null}. + * @see com.datastax.driver.core.Session */ - public Session getSession() { - return session; + public void setSession(Session session) { + Assert.notNull(session, "Session must not be null"); + this.session = session; } /** - * @param session The session to set, must not be{@literal null} + * Returns the Cassandra {@link Session} used by this template to perform Cassandra data access operations. + * + * @return the Cassandra {@link Session} used by this template. + * @see com.datastax.driver.core.Session */ - public void setSession(Session session) { - - Assert.notNull(session, "Session must not be null!"); - this.session = session; + public Session getSession() { + Assert.state(this.session != null, "Session was not properly initialized"); + return this.session; } } diff --git a/spring-cql/src/test/java/org/springframework/cassandra/core/CqlTemplateUnitTests.java b/spring-cql/src/test/java/org/springframework/cassandra/core/CqlTemplateUnitTests.java new file mode 100644 index 000000000..54da6bc12 --- /dev/null +++ b/spring-cql/src/test/java/org/springframework/cassandra/core/CqlTemplateUnitTests.java @@ -0,0 +1,267 @@ +/* + * Copyright 2013-2016 the original author or authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cassandra.core; + +import static org.hamcrest.Matchers.*; +import static org.junit.Assert.*; +import static org.mockito.Mockito.*; + +import java.util.Iterator; + +import org.junit.Before; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.ExpectedException; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.runners.MockitoJUnitRunner; +import org.springframework.dao.IncorrectResultSizeDataAccessException; + +import com.datastax.driver.core.ColumnDefinitions; +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.Row; +import com.datastax.driver.core.Session; + +/** + * The CqlTemplateUnitTests class is a test suite of test cases testing the contract and functionality of the + * {@link CqlTemplate} class. + * + * @author John Blum + * @see org.springframework.cassandra.core.CqlTemplate + * @since 1.5.0 + */ +// TODO: add many more unit tests until SUT test coverage is 100%! +@RunWith(MockitoJUnitRunner.class) +public class CqlTemplateUnitTests { + + @Rule public ExpectedException exception = ExpectedException.none(); + + private CqlTemplate template; + + @Mock private Session mockSession; + + @Before + public void setup() { + template = new CqlTemplate(mockSession); + } + + @Test + @SuppressWarnings("unchecked") + public void firstColumnToObjectReturnsColumnValue() { + final Row mockRow = mock(Row.class); + ColumnDefinitions mockColumnDefinitions = mock(ColumnDefinitions.class); + Iterator mockIterator = mock(Iterator.class); + final ColumnDefinitions.Definition mockColumnDefinition = mock(ColumnDefinitions.Definition.class); + + when(mockRow.getColumnDefinitions()).thenReturn(mockColumnDefinitions); + when(mockColumnDefinitions.iterator()).thenReturn(mockIterator); + when(mockIterator.hasNext()).thenReturn(true); + when(mockIterator.next()).thenReturn(mockColumnDefinition); + + template = new CqlTemplate() { + @Override + T columnToObject(Row row, ColumnDefinitions.Definition columnDefinition) { + assertThat(row, is(sameInstance(mockRow))); + assertThat(columnDefinition, is(sameInstance(mockColumnDefinition))); + return (T) "test"; + } + }; + + assertThat(String.valueOf(template.firstColumnToObject(mockRow)), is(equalTo("test"))); + + verify(mockRow, times(1)).getColumnDefinitions(); + verify(mockColumnDefinitions, times(1)).iterator(); + verify(mockIterator, times(1)).hasNext(); + verify(mockIterator, times(1)).next(); + verifyZeroInteractions(mockColumnDefinition); + } + + @Test + @SuppressWarnings("unchecked") + public void firstColumnToObjectReturnsNull() { + Row mockRow = mock(Row.class); + ColumnDefinitions mockColumnDefinitions = mock(ColumnDefinitions.class); + Iterator mockIterator = mock(Iterator.class); + + when(mockRow.getColumnDefinitions()).thenReturn(mockColumnDefinitions); + when(mockColumnDefinitions.iterator()).thenReturn(mockIterator); + when(mockIterator.hasNext()).thenReturn(false); + + assertThat(template.firstColumnToObject(mockRow), is(nullValue(Object.class))); + + verify(mockRow, times(1)).getColumnDefinitions(); + verify(mockColumnDefinitions, times(1)).iterator(); + verify(mockIterator, times(1)).hasNext(); + verify(mockIterator, never()).next(); + } + + @Test + @SuppressWarnings("unchecked") + public void processOneIsSuccessful() { + ResultSet mockResultSet = mock(ResultSet.class); + Row mockRow = mock(Row.class); + RowMapper mockRowMapper = mock(RowMapper.class); + + when(mockResultSet.one()).thenReturn(mockRow); + when(mockResultSet.isExhausted()).thenReturn(true); + when(mockRowMapper.mapRow(eq(mockRow), eq(0))).thenReturn("test"); + + assertThat(template.processOne(mockResultSet, mockRowMapper), is(equalTo("test"))); + + verify(mockResultSet, times(1)).one(); + verify(mockResultSet, times(1)).isExhausted(); + verify(mockRowMapper, times(1)).mapRow(eq(mockRow), eq(0)); + verifyZeroInteractions(mockRow); + } + + @Test + public void processOneThrowsIncorrectResultSetSizeDataAccessExceptionWhenNoRowsFound() { + ResultSet mockResultSet = mock(ResultSet.class); + RowMapper mockRowMapper = mock(RowMapper.class); + + when(mockResultSet.one()).thenReturn(null); + + try { + exception.expect(IncorrectResultSizeDataAccessException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage(containsString("expected 1, actual 0")); + + template.processOne(mockResultSet, mockRowMapper); + + } finally { + verify(mockResultSet, times(1)).one(); + verify(mockResultSet, never()).isExhausted(); + verifyZeroInteractions(mockRowMapper); + } + } + + @Test + public void processOneThrowsIncorrectResultSetSizeDataAccessExceptionWhenTooManyRowsFound() { + ResultSet mockResultSet = mock(ResultSet.class); + Row mockRow = mock(Row.class); + RowMapper mockRowMapper = mock(RowMapper.class); + + when(mockResultSet.one()).thenReturn(mockRow); + when(mockResultSet.isExhausted()).thenReturn(false); + + try { + exception.expect(IncorrectResultSizeDataAccessException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage("ResultSet size exceeds 1"); + + template.processOne(mockResultSet, mockRowMapper); + + } finally { + verify(mockResultSet, times(1)).one(); + verify(mockResultSet, times(1)).isExhausted(); + verifyZeroInteractions(mockRowMapper); + verifyZeroInteractions(mockRow); + } + } + + @Test + public void processOnePassingNullResultSetThrowsIllegalArgumentException() { + RowMapper mockRowMapper = mock(RowMapper.class); + + try { + exception.expect(IllegalArgumentException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage("ResultSet cannot be null"); + + template.processOne(null, mockRowMapper); + } finally { + verifyZeroInteractions(mockRowMapper); + } + } + + @Test + @SuppressWarnings("unchecked") + public void processOneWithRequiredTypeIsSuccessful() { + ResultSet mockResultSet = mock(ResultSet.class); + final Row mockRow = mock(Row.class); + + when(mockResultSet.one()).thenReturn(mockRow); + when(mockResultSet.isExhausted()).thenReturn(true); + + template = new CqlTemplate() { + @Override + protected Object firstColumnToObject(Row row) { + assertThat(row, is(equalTo(mockRow))); + return 1l; + } + }; + + Number value = template.processOne(mockResultSet, Long.class); + + assertThat(value, is(instanceOf(Long.class))); + assertThat(value.longValue(), is(equalTo(1l))); + + verify(mockResultSet, times(1)).one(); + verify(mockResultSet, times(1)).isExhausted(); + verifyZeroInteractions(mockRow); + } + + @Test + public void processOneWithRequiredTypeThrowsIncorrectResultSetSizeDataAccessExceptionWhenNoRowsFound() { + ResultSet mockResultSet = mock(ResultSet.class); + + when(mockResultSet.one()).thenReturn(null); + + try { + exception.expect(IncorrectResultSizeDataAccessException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage(containsString("expected 1, actual 0")); + + template.processOne(mockResultSet, Integer.class); + + } finally { + verify(mockResultSet, times(1)).one(); + verify(mockResultSet, never()).isExhausted(); + } + } + + @Test + public void processOneWithRequiredTypeThrowsIncorrectResultSetSizeDataAccessExceptionWhenTooManyRowsFound() { + ResultSet mockResultSet = mock(ResultSet.class); + Row mockRow = mock(Row.class); + + when(mockResultSet.one()).thenReturn(mockRow); + when(mockResultSet.isExhausted()).thenReturn(false); + + try { + exception.expect(IncorrectResultSizeDataAccessException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage(containsString("ResultSet size exceeds 1")); + + template.processOne(mockResultSet, Double.class); + + } finally { + verify(mockResultSet, times(1)).one(); + verify(mockResultSet, times(1)).isExhausted(); + verifyZeroInteractions(mockRow); + } + } + + @Test + public void processOneWithRequiredTypePassingNullResultSetThrowsIllegalArgumentException() { + exception.expect(IllegalArgumentException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage(is(equalTo("ResultSet cannot be null"))); + + template.processOne(null, String.class); + } +} diff --git a/spring-cql/src/test/java/org/springframework/cassandra/support/CassandraAccessorUnitTests.java b/spring-cql/src/test/java/org/springframework/cassandra/support/CassandraAccessorUnitTests.java new file mode 100644 index 000000000..8362c4938 --- /dev/null +++ b/spring-cql/src/test/java/org/springframework/cassandra/support/CassandraAccessorUnitTests.java @@ -0,0 +1,108 @@ +/* + * Copyright 2013-2016 the original author or authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cassandra.support; + +import static org.hamcrest.Matchers.*; +import static org.junit.Assert.*; + +import org.junit.Before; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.ExpectedException; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.runners.MockitoJUnitRunner; + +import com.datastax.driver.core.Session; + +/** + * The CassandraAccessorUnitTests class is a test suite of test cases testing the contract and functionality of the + * {@link CassandraAccessor} class. + * + * @author John Blum + * @see org.springframework.cassandra.support.CassandraAccessor + * @since 1.5.0 + */ +@RunWith(MockitoJUnitRunner.class) +public class CassandraAccessorUnitTests { + + private CassandraAccessor cassandraAccessor; + + @Mock private CassandraExceptionTranslator mockExceptionTranslator; + + @Rule public ExpectedException exception = ExpectedException.none(); + + @Mock private Session mockSession; + + @Before + public void setup() { + cassandraAccessor = new CassandraAccessor(); + } + + @Test + public void afterPropertiesSetWithUnitializedSessionThrowsIllegalStateException() { + exception.expect(IllegalStateException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage("Session must not be null"); + + cassandraAccessor.afterPropertiesSet(); + } + + @Test + public void setAndGetExceptionTranslator() { + cassandraAccessor.setExceptionTranslator(mockExceptionTranslator); + assertThat(cassandraAccessor.getExceptionTranslator(), is(sameInstance(mockExceptionTranslator))); + } + + @Test + public void setExceptionTranslatorToNullThrowsIllegalArgumentException() { + exception.expect(IllegalArgumentException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage(is(equalTo("CassandraExceptionTranslator must not be null"))); + + cassandraAccessor.setExceptionTranslator(null); + } + + @Test + public void getUninitializedExceptionTranslatorReturnsDefault() { + assertThat(cassandraAccessor.getExceptionTranslator(), is(equalTo(cassandraAccessor.exceptionTranslator))); + } + + @Test + public void setAndGetSession() { + cassandraAccessor.setSession(mockSession); + assertThat(cassandraAccessor.getSession(), is(sameInstance(mockSession))); + } + + @Test + public void setSessionToNullThrowsIllegalArgumentException() { + exception.expect(IllegalArgumentException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage(is(equalTo("Session must not be null"))); + + cassandraAccessor.setSession(null); + } + + @Test + public void getUninitializedSessionThrowsIllegalStateException() { + exception.expect(IllegalStateException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage(is(equalTo("Session was not properly initialized"))); + + cassandraAccessor.getSession(); + } +} diff --git a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/core/CqlOperationsIntegrationTests.java b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/core/CqlOperationsIntegrationTests.java index 14d3b8f65..08249d5a8 100644 --- a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/core/CqlOperationsIntegrationTests.java +++ b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/core/CqlOperationsIntegrationTests.java @@ -49,6 +49,7 @@ import org.springframework.cassandra.core.WriteOptions; import org.springframework.cassandra.core.keyspace.CreateTableSpecification; import org.springframework.cassandra.test.integration.AbstractKeyspaceCreatingIntegrationTest; import org.springframework.dao.DataAccessException; +import org.springframework.dao.IncorrectResultSizeDataAccessException; import org.springframework.util.CollectionUtils; import com.datastax.driver.core.BoundStatement; @@ -607,19 +608,18 @@ public class CqlOperationsIntegrationTests extends AbstractKeyspaceCreatingInteg /** * Test that CQL for QueryForObject must only return 1 row or an IllegalArgumentException is thrown. */ - @Test(expected = IllegalArgumentException.class) + @Test(expected = IncorrectResultSizeDataAccessException.class) public void queryForObjectTestCqlStringRowMapperNotOneRowReturned() { // Insert our 3 test books. insertTestObjectArray(); @SuppressWarnings("unused") - Book book = cqlTemplate.queryForObject("select * from book where isbn in ('1234','2345','3456')", + Book book = cqlTemplate.queryForObject("SELECT * FROM book WHERE isbn IN('1234','2345','3456')", new RowMapper() { @Override public Book mapRow(Row row, int rowNum) throws DriverException { - Book b = rowToBook(row); - return b; + return rowToBook(row); } }); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java index 4fd03b125..3610d929e 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java @@ -17,8 +17,8 @@ package org.springframework.data.cassandra.core; import java.util.List; -import org.springframework.cassandra.core.CqlOperations; import org.springframework.cassandra.core.Cancellable; +import org.springframework.cassandra.core.CqlOperations; import org.springframework.cassandra.core.QueryForObjectListener; import org.springframework.cassandra.core.QueryOptions; import org.springframework.cassandra.core.WriteOptions; diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java index 40dfa9843..4af90b4c3 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java @@ -97,11 +97,16 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation */ public CassandraTemplate(Session session, CassandraConverter converter) { setSession(session); - setConverter(converter != null ? converter : getDefaultCassandraConverter()); + setConverter(resolveConverter(converter)); } - private static CassandraConverter getDefaultCassandraConverter() { + /* (non-Javadoc) */ + private static CassandraConverter resolveConverter(CassandraConverter cassandraConverter) { + return (cassandraConverter != null ? cassandraConverter : getDefaultCassandraConverter()); + } + /* (non-Javadoc) */ + private static CassandraConverter getDefaultCassandraConverter() { MappingCassandraConverter mappingCassandraConverter = new MappingCassandraConverter(); mappingCassandraConverter.afterPropertiesSet(); return mappingCassandraConverter; @@ -110,8 +115,8 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation /** * Set the {@link CassandraConverter} used by this template to perform conversions. * - * @param cassandraConverter Converter used to perform conversion of Cassandra data types to entity types. - * Must not be {@literal null}. + * @param cassandraConverter Converter used to perform conversion of Cassandra data types to entity types. Must not be + * {@literal null}. */ public void setConverter(CassandraConverter cassandraConverter) { @@ -146,8 +151,8 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation super.afterPropertiesSet(); - Assert.notNull(cassandraConverter, "CassandraConverter must not be null!"); - Assert.notNull(mappingContext, "CassandraMappingContext must not be null!"); + Assert.notNull(cassandraConverter, "CassandraConverter must not be null"); + Assert.notNull(mappingContext, "CassandraMappingContext must not be null"); } @Override