From ed46d742eb066de2a746266d80d127847496a9d0 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Thu, 3 Dec 2020 17:07:15 +0100 Subject: [PATCH] DATACASS-510 - Introduce support for prepared statements using CassandraTemplate. We now support prepared statement usage through CassandraTemplate and its asynchronous and reactive variants. All statements created or received by CassandraTemplate will be prepared. CassandraTemplate is an infrastructure class for repositories so prepared statements will affect repositories, too. --- .../core/AsyncCassandraOperations.java | 12 ++ .../core/AsyncCassandraTemplate.java | 189 +++++++++++++---- .../cassandra/core/CassandraOperations.java | 26 ++- .../cassandra/core/CassandraTemplate.java | 194 +++++++++++++---- .../core/PreparedStatementDelegate.java | 120 +++++++++++ .../core/ReactiveCassandraOperations.java | 26 ++- .../core/ReactiveCassandraTemplate.java | 199 +++++++++++++----- .../cql/AsyncPreparedStatementCreator.java | 11 +- .../core/cql/PreparedStatementCreator.java | 5 + .../cql/ReactivePreparedStatementCreator.java | 13 +- .../query/AbstractCassandraQuery.java | 2 +- .../query/AbstractReactiveCassandraQuery.java | 5 +- .../query/CassandraQueryExecution.java | 33 ++- .../ReactiveCassandraQueryExecution.java | 10 +- ...syncCassandraTemplateIntegrationTests.java | 10 + ...atePreparedStatementsIntegrationTests.java | 30 +++ .../core/AsyncCassandraTemplateUnitTests.java | 1 + .../CassandraTemplateIntegrationTests.java | 11 + ...latePreparedStatementIntegrationTests.java | 30 +++ .../core/CassandraTemplateUnitTests.java | 1 + ...tiveCassandraTemplateIntegrationTests.java | 10 + ...latePreparedStatementIntegrationTests.java | 30 +++ .../ReactiveCassandraTemplateUnitTests.java | 1 + .../example/CassandraTemplateExamples.java | 49 +++++ .../example/CqlTemplateExamples.java | 7 + ...yMethodParameterTypesIntegrationTests.java | 8 + .../CoroutineRepositoryUnitTests.kt | 3 +- src/main/asciidoc/new-features.adoc | 5 + .../reference/cassandra-repositories.adoc | 1 + src/main/asciidoc/reference/cassandra.adoc | 77 ++++++- .../reactive-cassandra-repositories.adoc | 2 + .../reference/reactive-cassandra.adoc | 2 +- 32 files changed, 934 insertions(+), 189 deletions(-) create mode 100644 spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/PreparedStatementDelegate.java create mode 100644 spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplatePreparedStatementsIntegrationTests.java create mode 100644 spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplatePreparedStatementIntegrationTests.java create mode 100644 spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplatePreparedStatementIntegrationTests.java create mode 100644 spring-data-cassandra/src/test/java/org/springframework/data/cassandra/example/CassandraTemplateExamples.java diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraOperations.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraOperations.java index d170a1cf3..fe0fa8024 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraOperations.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraOperations.java @@ -29,6 +29,7 @@ import org.springframework.data.cassandra.core.query.Update; import org.springframework.data.domain.Slice; import org.springframework.util.concurrent.ListenableFuture; +import com.datastax.oss.driver.api.core.cql.AsyncResultSet; import com.datastax.oss.driver.api.core.cql.Statement; /** @@ -102,6 +103,17 @@ public interface AsyncCassandraOperations { // Methods dealing with com.datastax.oss.driver.api.core.cql.Statement // ------------------------------------------------------------------------- + /** + * Execute the a Cassandra {@link Statement}. Any errors that result from executing this command will be converted + * into Spring's DAO exception hierarchy. + * + * @param statement a Cassandra {@link Statement}, must not be {@literal null}. + * @return the {@link AsyncResultSet}. + * @throws DataAccessException if there is any problem executing the query. + * @since 3.2 + */ + ListenableFuture execute(Statement statement) throws DataAccessException; + /** * Execute a {@code SELECT} query and convert the resulting items to a {@link List} of entities. * diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java index e550fb4c4..bed8aa118 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java @@ -22,6 +22,9 @@ import java.util.function.Function; import java.util.stream.Collectors; import java.util.stream.StreamSupport; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import org.springframework.beans.BeansException; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; @@ -29,18 +32,12 @@ import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisherAware; import org.springframework.dao.DataAccessException; import org.springframework.dao.OptimisticLockingFailureException; +import org.springframework.dao.support.DataAccessUtils; import org.springframework.data.cassandra.SessionFactory; import org.springframework.data.cassandra.core.EntityOperations.AdaptibleEntity; import org.springframework.data.cassandra.core.convert.CassandraConverter; import org.springframework.data.cassandra.core.convert.MappingCassandraConverter; -import org.springframework.data.cassandra.core.cql.AsyncCqlOperations; -import org.springframework.data.cassandra.core.cql.AsyncCqlTemplate; -import org.springframework.data.cassandra.core.cql.AsyncSessionCallback; -import org.springframework.data.cassandra.core.cql.CassandraAccessor; -import org.springframework.data.cassandra.core.cql.CqlExceptionTranslator; -import org.springframework.data.cassandra.core.cql.CqlProvider; -import org.springframework.data.cassandra.core.cql.QueryOptions; -import org.springframework.data.cassandra.core.cql.WriteOptions; +import org.springframework.data.cassandra.core.cql.*; import org.springframework.data.cassandra.core.cql.session.DefaultSessionFactory; import org.springframework.data.cassandra.core.cql.util.CassandraFutureAdapter; import org.springframework.data.cassandra.core.cql.util.StatementBuilder; @@ -59,6 +56,7 @@ import org.springframework.data.domain.Slice; import org.springframework.data.mapping.callback.EntityCallbacks; import org.springframework.data.projection.ProjectionFactory; import org.springframework.data.projection.SpelAwareProxyProjectionFactory; +import org.springframework.data.util.Streamable; import org.springframework.lang.Nullable; import org.springframework.scheduling.annotation.AsyncResult; import org.springframework.util.Assert; @@ -69,6 +67,9 @@ import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.DriverException; import com.datastax.oss.driver.api.core.config.DefaultDriverOption; import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.BoundStatement; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.cql.ResultSet; import com.datastax.oss.driver.api.core.cql.Row; import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; @@ -89,6 +90,13 @@ import com.datastax.oss.driver.api.querybuilder.update.Update; * Can be used within a service implementation via direct instantiation with a {@link CqlSession} reference, or get * prepared in an application context and given to services as bean reference. *

+ * This class supports the use of prepared statements when enabling {@link #setUsePreparedStatements(boolean)}. All + * statements created by methods of this class (such as {@link #select(Query, Class)} or + * {@link #update(Query, org.springframework.data.cassandra.core.query.Update, Class)} will be executed as prepared + * statements. Also, statements accepted by methods (such as {@link #select(String, Class)} or + * {@link #select(Statement, Class) and others}) will be prepared prior to execution. Note that {@link Statement} + * objects passed to methods must be {@link SimpleStatement} so that these can be prepared. + *

* Note: The {@link CqlSession} should always be configured as a bean in the application context, in the first case * given to the service directly, in the second case to the prepared template. * @@ -100,6 +108,8 @@ import com.datastax.oss.driver.api.querybuilder.update.Update; public class AsyncCassandraTemplate implements AsyncCassandraOperations, ApplicationEventPublisherAware, ApplicationContextAware { + private final Logger logger = LoggerFactory.getLogger(getClass()); + private final AsyncCqlOperations cqlOperations; private final CassandraConverter converter; @@ -116,6 +126,8 @@ public class AsyncCassandraTemplate private @Nullable EntityCallbacks entityCallbacks; + private boolean usePreparedStatements = true; + /** * Creates an instance of {@link AsyncCassandraTemplate} initialized with the given {@link CqlSession} and a default * {@link MappingCassandraConverter}. @@ -136,7 +148,7 @@ public class AsyncCassandraTemplate * @param converter {@link CassandraConverter} used to convert between Java and Cassandra types; must not be * {@literal null}. * @see CassandraConverter - * @see Session + * @see CqlSession */ public AsyncCassandraTemplate(CqlSession session, CassandraConverter converter) { this(new DefaultSessionFactory(session), converter); @@ -150,7 +162,7 @@ public class AsyncCassandraTemplate * @param converter {@link CassandraConverter} used to convert between Java and Cassandra types; must not be * {@literal null}. * @see CassandraConverter - * @see Session + * @see CqlSession */ public AsyncCassandraTemplate(SessionFactory sessionFactory, CassandraConverter converter) { this(new AsyncCqlTemplate(sessionFactory), converter); @@ -164,7 +176,7 @@ public class AsyncCassandraTemplate * @param converter {@link CassandraConverter} used to convert between Java and Cassandra types; must not be * {@literal null}. * @see CassandraConverter - * @see Session + * @see CqlSession */ public AsyncCassandraTemplate(AsyncCqlTemplate asyncCqlTemplate, CassandraConverter converter) { @@ -226,6 +238,32 @@ public class AsyncCassandraTemplate return this.converter; } + /** + * Returns whether this instance is configured to use {@link PreparedStatement prepared statements}. If enabled + * (default), then all persistence methods (such as {@link #select}, {@link #update}, and others) will make use of + * prepared statements. Note that methods accepting a {@link Statement} must be called with {@link SimpleStatement} + * instances to participate in statement preparation. + * + * @return {@literal true} if prepared statements usage is enabled; {@literal false} otherwise. + * @since 3.2 + */ + public boolean isUsePreparedStatements() { + return usePreparedStatements; + } + + /** + * Enable/disable {@link PreparedStatement prepared statements} usage. If enabled (default), then all persistence + * methods (such as {@link #select}, {@link #update}, and others) will make use of prepared statements. Note that + * methods accepting a {@link Statement} must be called with {@link SimpleStatement} instances to participate in + * statement preparation. + * + * @param usePreparedStatements whether to use prepared statements. + * @since 3.2 + */ + public void setUsePreparedStatements(boolean usePreparedStatements) { + this.usePreparedStatements = usePreparedStatements; + } + /** * Returns the {@link EntityOperations} used to perform data access operations on an entity inside a Cassandra data * source. @@ -314,6 +352,17 @@ public class AsyncCassandraTemplate // Methods dealing with com.datastax.oss.driver.api.core.cql.Statement // ------------------------------------------------------------------------- + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.AsyncCassandraOperations#select(com.datastax.oss.driver.api.core.cql.Statement, java.lang.Class) + */ + @Override + public ListenableFuture execute(Statement statement) throws DataAccessException { + + Assert.notNull(statement, "Statement must not be null"); + + return doQueryForResultSet(statement); + } + /* (non-Javadoc) * @see org.springframework.data.cassandra.core.AsyncCassandraOperations#select(com.datastax.oss.driver.api.core.cql.Statement, java.lang.Class) */ @@ -325,7 +374,7 @@ public class AsyncCassandraTemplate Function mapper = getMapper(entityClass, entityClass, EntityQueryUtils.getTableName(statement)); - return getAsyncCqlOperations().query(statement, (row, rowNum) -> mapper.apply(row)); + return doQuery(statement, (row, rowNum) -> mapper.apply(row)); } /* (non-Javadoc) @@ -341,7 +390,7 @@ public class AsyncCassandraTemplate Function mapper = getMapper(entityClass, entityClass, EntityQueryUtils.getTableName(statement)); - return getAsyncCqlOperations().query(statement, row -> { + return doQuery(statement, row -> { entityConsumer.accept(mapper.apply(row)); }); } @@ -364,7 +413,7 @@ public class AsyncCassandraTemplate Assert.notNull(statement, "Statement must not be null"); Assert.notNull(entityClass, "Entity type must not be null"); - ListenableFuture resultSet = getAsyncCqlOperations().queryForResultSet(statement); + ListenableFuture resultSet = doQueryForResultSet(statement); Function mapper = getMapper(entityClass, entityClass, EntityQueryUtils.getTableName(statement)); @@ -439,8 +488,8 @@ public class AsyncCassandraTemplate Assert.notNull(update, "Update must not be null"); Assert.notNull(entityClass, "Entity type must not be null"); - return getAsyncCqlOperations() - .execute(getStatementFactory().update(query, update, getRequiredPersistentEntity(entityClass)).build()); + return doExecute(getStatementFactory().update(query, update, getRequiredPersistentEntity(entityClass)).build(), + AsyncResultSet::wasApplied); } /* (non-Javadoc) @@ -463,7 +512,7 @@ public class AsyncCassandraTemplate maybeEmitEvent(new BeforeDeleteEvent<>(delete, entityClass, tableName)); - ListenableFuture future = getAsyncCqlOperations().execute(delete); + ListenableFuture future = doExecute(delete, AsyncResultSet::wasApplied); future.addCallback(success -> maybeEmitEvent(new AfterDeleteEvent<>(delete, entityClass, tableName)), e -> {}); @@ -504,7 +553,13 @@ public class AsyncCassandraTemplate SimpleStatement statement = countStatement.build(); - ListenableFuture result = getAsyncCqlOperations().queryForObject(statement, Long.class); + ListenableFuture result = doExecute(statement, it -> { + + SingleColumnRowMapper mapper = SingleColumnRowMapper.newInstance(Long.class); + + Row row = DataAccessUtils.requiredSingleResult(Streamable.of(it.currentPage()).toList()); + return mapper.mapRow(row, 0); + }); return new MappingListenableFutureAdapter<>(result, it -> it != null ? it : 0L); } @@ -523,8 +578,7 @@ public class AsyncCassandraTemplate StatementBuilder select = getStatementFactory() .selectOneById(id, entity, entity.getTableName()); - return new MappingListenableFutureAdapter<>(getAsyncCqlOperations().queryForResultSet(select.build()), - resultSet -> resultSet.one() != null); + return doExecute(select.build(), resultSet -> resultSet.one() != null); } /* (non-Javadoc) @@ -539,8 +593,7 @@ public class AsyncCassandraTemplate StatementBuilder select = getStatementFactory() .select(query.limit(1), getRequiredPersistentEntity(entityClass), getTableName(entityClass)); - return new MappingListenableFutureAdapter<>(getAsyncCqlOperations().queryForResultSet(select.build()), - resultSet -> resultSet.one() != null); + return doExecute(select.build(), resultSet -> resultSet.one() != null); } /* (non-Javadoc) @@ -557,8 +610,7 @@ public class AsyncCassandraTemplate StatementBuilder countStatement = getStatementFactory().count(query, getRequiredPersistentEntity(entityClass), tableName); - SimpleStatement statement = countStatement.build(); - Long count = getCqlOperations().queryForObject(statement, Long.class); - - return count != null ? count : 0L; + return doQueryForObject(countStatement.build(), Long.class); } /* (non-Javadoc) @@ -561,7 +610,7 @@ public class CassandraTemplate implements CassandraOperations, ApplicationEventP CassandraPersistentEntity entity = getRequiredPersistentEntity(entityClass); StatementBuilder select = getStatementFactory().select(query.limit(1), getRequiredPersistentEntity(entityClass), tableName); - return getCqlOperations().queryForResultSet(select.build()).one() != null; + return doQueryForResultSet(select.build()).one() != null; } /* (non-Javadoc) @@ -597,7 +646,7 @@ public class CassandraTemplate implements CassandraOperations, ApplicationEventP CqlIdentifier tableName = entity.getTableName(); StatementBuilder count = getStatementFactory().count(query, getRequiredPersistentEntity(entityClass), tableName); - return getReactiveCqlOperations().queryForObject(count.build(), Long.class).switchIfEmpty(Mono.just(0L)); + SingleColumnRowMapper mapper = SingleColumnRowMapper.newInstance(Long.class); + + Mono mono = doExecuteAndFlatMap(count.build(), rs -> rs.rows() // + .map(it -> mapper.mapRow(it, 0)) // + .buffer() // + .map(DataAccessUtils::requiredSingleResult).next()); + + return mono.switchIfEmpty(Mono.just(0L)); } /* (non-Javadoc) @@ -506,7 +558,7 @@ public class ReactiveCassandraTemplate CassandraPersistentEntity entity = getRequiredPersistentEntity(entityClass); StatementBuilder builder = getStatementFactory().select(query.limit(1), getRequiredPersistentEntity(entityClass), tableName); - return getReactiveCqlOperations().queryForRows(builder.build()).hasElements(); + return doQuery(builder.build(), (row, rowNum) -> row).hasElements(); } /* (non-Javadoc) @@ -734,7 +786,7 @@ public class ReactiveCassandraTemplate StatementBuilder builder = getStatementFactory().deleteById(id, entity, tableName); SimpleStatement delete = builder.build(); - Mono result = getReactiveCqlOperations().execute(delete) + Mono result = doExecute(delete, ReactiveResultSet::wasApplied) .doOnSubscribe(it -> maybeEmitEvent(new BeforeDeleteEvent<>(delete, entityClass, tableName))); return result.doOnNext(it -> maybeEmitEvent(new AfterDeleteEvent<>(delete, entityClass, tableName))); @@ -752,7 +804,7 @@ public class ReactiveCassandraTemplate Truncate truncate = QueryBuilder.truncate(tableName); SimpleStatement statement = truncate.build(); - Mono result = getReactiveCqlOperations().execute(statement) + Mono result = doExecute(statement, ReactiveResultSet::wasApplied) .doOnSubscribe(it -> maybeEmitEvent(new BeforeDeleteEvent<>(statement, entityClass, tableName))); return result.doOnNext(it -> maybeEmitEvent(new AfterDeleteEvent<>(statement, entityClass, tableName))).then(); @@ -810,13 +862,12 @@ public class ReactiveCassandraTemplate maybeEmitEvent(new BeforeSaveEvent<>(entity, tableName, statement)); return maybeCallBeforeSave(entity, tableName, statement).flatMapMany(entityToSave -> { - Flux execute = getReactiveCqlOperations().execute(new StatementCallback(statement)); + Mono execute = doExecuteAndFlatMap(statement, ReactiveCassandraTemplate::toWriteResult); return execute.map(it -> EntityWriteResult.of(it, entityToSave)).handle(handler) // .doOnNext(it -> maybeEmitEvent(new AfterSaveEvent<>(entityToSave, tableName))); }).next(); }); - } private Mono executeDelete(Object entity, CqlIdentifier tableName, SimpleStatement statement, @@ -824,12 +875,46 @@ public class ReactiveCassandraTemplate maybeEmitEvent(new BeforeDeleteEvent<>(statement, entity.getClass(), tableName)); - Flux execute = getReactiveCqlOperations().execute(new StatementCallback(statement)); + Mono execute = doExecuteAndFlatMap(statement, ReactiveCassandraTemplate::toWriteResult); return execute.map(it -> EntityWriteResult.of(it, entity)).handle(handler) // .doOnSubscribe(it -> maybeEmitEvent(new BeforeSaveEvent<>(entity, tableName, statement))) // - .doOnNext(it -> maybeEmitEvent(new AfterDeleteEvent<>(statement, entity.getClass(), tableName))) // - .next(); + .doOnNext(it -> maybeEmitEvent(new AfterDeleteEvent<>(statement, entity.getClass(), tableName))); + } + + private Flux doQuery(Statement statement, RowMapper rowMapper) { + + if (PreparedStatementDelegate.canPrepare(isUsePreparedStatements(), statement, logger)) { + + PreparedStatementHandler statementHandler = new PreparedStatementHandler(statement); + return getReactiveCqlOperations().query(statementHandler, statementHandler, rowMapper); + } + + return getReactiveCqlOperations().query(statement, rowMapper); + } + + private Mono doExecute(Statement statement, Function mappingFunction) { + + if (PreparedStatementDelegate.canPrepare(isUsePreparedStatements(), statement, logger)) { + + PreparedStatementHandler statementHandler = new PreparedStatementHandler(statement); + return getReactiveCqlOperations() + .query(statementHandler, statementHandler, rs -> Mono.just(mappingFunction.apply(rs))).next(); + } + + return getReactiveCqlOperations().queryForResultSet(statement).map(mappingFunction); + } + + private Mono doExecuteAndFlatMap(Statement statement, + Function> mappingFunction) { + + if (PreparedStatementDelegate.canPrepare(isUsePreparedStatements(), statement, logger)) { + + PreparedStatementHandler statementHandler = new PreparedStatementHandler(statement); + return getReactiveCqlOperations().query(statementHandler, statementHandler, mappingFunction::apply).next(); + } + + return getReactiveCqlOperations().queryForResultSet(statement).flatMap(mappingFunction); } private int getConfiguredPageSize(DriverContext context) { @@ -875,6 +960,11 @@ public class ReactiveCassandraTemplate }; } + static Mono toWriteResult(ReactiveResultSet resultSet) { + return resultSet.rows().collectList() + .map(rows -> new WriteResult(resultSet.getAllExecutionInfo(), resultSet.wasApplied(), rows)); + } + private Class resolveTypeToRead(Class entityType, Class targetType) { return targetType.isInterface() || targetType.isAssignableFrom(entityType) ? entityType : targetType; } @@ -913,21 +1003,37 @@ public class ReactiveCassandraTemplate return Mono.just(object); } - static class StatementCallback implements ReactiveSessionCallback, CqlProvider { + /** + * Utility class to prepare a {@link SimpleStatement} and bind values associated with the statement to a + * {@link BoundStatement}. + * + * @since 3.2 + */ + private static class PreparedStatementHandler + implements ReactivePreparedStatementCreator, PreparedStatementBinder, CqlProvider { private final SimpleStatement statement; - StatementCallback(SimpleStatement statement) { - this.statement = statement; + public PreparedStatementHandler(Statement statement) { + this.statement = PreparedStatementDelegate.getStatementForPrepare(statement); } /* * (non-Javadoc) - * @see org.springframework.data.cassandra.core.cql.ReactiveSessionCallback#doInSession(org.springframework.data.cassandra.ReactiveSession) + * @see org.springframework.data.cassandra.core.cql.ReactivePreparedStatementCreator#doInSession(org.springframework.data.cassandra.ReactiveSession) */ @Override - public Publisher doInSession(ReactiveSession session) throws DriverException, DataAccessException { - return session.execute(this.statement).flatMap(StatementCallback::toWriteResult); + public Mono createPreparedStatement(ReactiveSession session) throws DriverException { + return session.prepare(statement); + } + + /* + * (non-Javadoc) + * @see org.springframework.data.cassandra.core.cql.PreparedStatementBinder#bindValues(com.datastax.oss.driver.api.core.cql.PreparedStatement) + */ + @Override + public BoundStatement bindValues(PreparedStatement ps) throws DriverException { + return PreparedStatementDelegate.bind(statement, ps); } /* @@ -936,12 +1042,7 @@ public class ReactiveCassandraTemplate */ @Override public String getCql() { - return this.statement.getQuery(); - } - - private static Mono toWriteResult(ReactiveResultSet resultSet) { - return resultSet.rows().collectList() - .map(rows -> new WriteResult(resultSet.getAllExecutionInfo(), resultSet.wasApplied(), rows)); + return statement.getQuery(); } } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/AsyncPreparedStatementCreator.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/AsyncPreparedStatementCreator.java index e917a318b..90ac0cfdc 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/AsyncPreparedStatementCreator.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/AsyncPreparedStatementCreator.java @@ -15,12 +15,12 @@ */ package org.springframework.data.cassandra.core.cql; +import org.springframework.util.concurrent.ListenableFuture; + import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.DriverException; import com.datastax.oss.driver.api.core.cql.PreparedStatement; -import org.springframework.util.concurrent.ListenableFuture; - /** * One of the two central callback interfaces used by the {@link AsyncCqlTemplate} class. This interface prepares a CQL * statement returning a {@link org.springframework.util.concurrent.ListenableFuture} given a {@link CqlSession}, @@ -30,13 +30,14 @@ import org.springframework.util.concurrent.ListenableFuture; * concern themselves with {@link DriverException}s that may be thrown from operations they attempt. The * {@link AsyncCqlTemplate} class will catch and handle {@link DriverException}s appropriately. *

- * A {@link AsyncPreparedStatementCreator} should also implement the {@link CqlProvider} interface if it is able to - * provide the CQL it uses for {@link PreparedStatement} creation. This allows for better contextual information in case - * of exceptions. + * Classes implementing this interface should also implement the {@link CqlProvider} interface if it is able to provide + * the CQL it uses for {@link PreparedStatement} creation. This allows for better contextual information in case of + * exceptions. * * @author Mark Paluch * @since 2.0 * @see AsyncCqlTemplate#execute(AsyncPreparedStatementCreator, PreparedStatementCallback) + * @see CqlProvider */ @FunctionalInterface public interface AsyncPreparedStatementCreator { diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/PreparedStatementCreator.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/PreparedStatementCreator.java index c12678df6..bb4084fee 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/PreparedStatementCreator.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/PreparedStatementCreator.java @@ -26,11 +26,16 @@ import com.datastax.oss.driver.api.core.cql.PreparedStatement; *

* Implementations do not need to concern themselves with {@link DriverException}s that may be thrown from * operations they attempt. The {@link CqlTemplate} class will catch and handle {@link DriverException}s appropriately. + *

+ * Classes implementing this interface should also implement the {@link CqlProvider} interface if it is able to provide + * the CQL it uses for {@link PreparedStatement} creation. This allows for better contextual information in case of + * exceptions. * * @author David Webb * @author Mark Paluch * @see CqlTemplate#execute(PreparedStatementCreator, PreparedStatementCallback) * @see CqlTemplate#query(PreparedStatementCreator, RowCallbackHandler) + * @see CqlProvider */ @FunctionalInterface public interface PreparedStatementCreator { diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactivePreparedStatementCreator.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactivePreparedStatementCreator.java index 785266da8..b0cdac7d5 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactivePreparedStatementCreator.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactivePreparedStatementCreator.java @@ -15,12 +15,13 @@ */ package org.springframework.data.cassandra.core.cql; -import com.datastax.oss.driver.api.core.DriverException; -import com.datastax.oss.driver.api.core.cql.PreparedStatement; import reactor.core.publisher.Mono; import org.springframework.data.cassandra.ReactiveSession; +import com.datastax.oss.driver.api.core.DriverException; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; + /** * One of the two central callback interfaces used by the {@link ReactiveCqlTemplate} class. This interface creates a * {@link PreparedStatement} given a {@link ReactiveSession}, provided by the {@link ReactiveCqlTemplate} class. @@ -29,12 +30,14 @@ import org.springframework.data.cassandra.ReactiveSession; * concern themselves with {@link DriverException}s that may be thrown from operations they attempt. The * {@link ReactiveCqlTemplate} class will catch and handle {@link DriverException}s appropriately. *

- * A {@link ReactivePreparedStatementCreator} should also implement the {@link CqlProvider} interface if it is able to - * provide the CQL it uses for {@link PreparedStatement} creation. This allows for better contextual information in case - * of exceptions. + * Classes implementing this interface should also implement the {@link CqlProvider} interface if it is able to provide + * the CQL it uses for {@link PreparedStatement} creation. This allows for better contextual information in case of + * exceptions. * * @author Mark Paluch * @since 2.0 + * @see ReactiveCqlTemplate#execute(ReactivePreparedStatementCreator, ReactivePreparedStatementCallback) + * @see CqlProvider */ @FunctionalInterface public interface ReactivePreparedStatementCreator { diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractCassandraQuery.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractCassandraQuery.java index 82540a67b..8b5ba5947 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractCassandraQuery.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractCassandraQuery.java @@ -149,7 +149,7 @@ public abstract class AbstractCassandraQuery extends CassandraRepositoryQuerySup } else if (isExistsQuery()) { return new ExistsExecution(getOperations()); } else if (isModifyingQuery()) { - return ((statement, type) -> getOperations().getCqlOperations().queryForResultSet(statement).wasApplied()); + return ((statement, type) -> getOperations().execute(statement).wasApplied()); } else { return new SingleEntityExecution(getOperations(), isLimiting()); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java index afd31ffe6..299e5284e 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractReactiveCassandraQuery.java @@ -146,8 +146,9 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository } else if (isExistsQuery()) { return new ExistsExecution(getReactiveCassandraOperations()); } else if (isModifyingQuery()) { - return (statement, type) -> getReactiveCassandraOperations().getReactiveCqlOperations() - .queryForResultSet(statement).map(ReactiveResultSet::wasApplied); + + return (statement, type) -> getReactiveCassandraOperations().execute(statement) + .map(ReactiveResultSet::wasApplied); } else { return new SingleEntityExecution(getReactiveCassandraOperations(), isLimiting()); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/CassandraQueryExecution.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/CassandraQueryExecution.java index c626e7dc2..400dfd113 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/CassandraQueryExecution.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/CassandraQueryExecution.java @@ -15,7 +15,6 @@ */ package org.springframework.data.cassandra.repository.query; -import java.util.Iterator; import java.util.List; import org.springframework.core.convert.converter.Converter; @@ -34,7 +33,6 @@ import org.springframework.data.repository.query.ReturnedType; import org.springframework.lang.Nullable; import org.springframework.util.ClassUtils; -import com.datastax.oss.driver.api.core.cql.ResultSet; import com.datastax.oss.driver.api.core.cql.Row; import com.datastax.oss.driver.api.core.cql.Statement; @@ -193,25 +191,22 @@ interface CassandraQueryExecution { @Override public Object execute(Statement statement, Class type) { - ResultSet resultSet = this.operations.getCqlOperations().queryForResultSet(statement); + List resultSet = this.operations.select(statement, Row.class); - Iterator iterator = resultSet.iterator(); - - if (iterator.hasNext()) { - - Row row = iterator.next(); - - if (!iterator.hasNext() && ProjectionUtil.qualifiesAsCountProjection(row)) { - - Object object = row.getObject(0); - - return ((Number) object).longValue() > 0; - } - - return true; + if (resultSet.isEmpty()) { + return false; } - return false; + Row row = resultSet.get(0); + + if (resultSet.size() == 1 && ProjectionUtil.qualifiesAsCountProjection(row)) { + + Object object = row.getObject(0); + + return ((Number) object).longValue() > 0; + } + + return true; } } @@ -233,7 +228,7 @@ interface CassandraQueryExecution { */ @Override public Object execute(Statement statement, Class type) { - return operations.getCqlOperations().queryForResultSet(statement); + return operations.execute(statement); } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraQueryExecution.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraQueryExecution.java index 3fc1d59d3..22daea3b0 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraQueryExecution.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraQueryExecution.java @@ -142,17 +142,17 @@ interface ReactiveCassandraQueryExecution { @Override public Publisher execute(Statement statement, Class type) { - return operations.select(statement, type).buffer(2).map(objects -> { + return operations.select(statement, type).buffer(2).handle((objects, sink) -> { if (objects.isEmpty()) { - return null; + return; } if (objects.size() == 1 || limiting) { - return objects.get(0); + sink.next(objects.get(0)); } - throw new IncorrectResultSizeDataAccessException(1, objects.size()); + sink.error(new IncorrectResultSizeDataAccessException(1, objects.size())); }); } } @@ -178,7 +178,7 @@ interface ReactiveCassandraQueryExecution { @Override public Publisher execute(Statement statement, Class type) { - Mono> rows = this.operations.getReactiveCqlOperations().queryForRows(statement).buffer(2).next(); + Mono> rows = this.operations.select(statement, Row.class).buffer(2).next(); return rows.map(it -> { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateIntegrationTests.java index 8370548a9..a27d7e5e3 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateIntegrationTests.java @@ -57,6 +57,7 @@ class AsyncCassandraTemplateIntegrationTests extends AbstractKeyspaceCreatingInt MappingCassandraConverter converter = new MappingCassandraConverter(); CassandraTemplate cassandraTemplate = new CassandraTemplate(session, converter); template = new AsyncCassandraTemplate(new AsyncCqlTemplate(session), converter); + prepareTemplate(template); SchemaTestUtils.potentiallyCreateTableFor(User.class, cassandraTemplate); SchemaTestUtils.potentiallyCreateTableFor(UserToken.class, cassandraTemplate); @@ -64,6 +65,15 @@ class AsyncCassandraTemplateIntegrationTests extends AbstractKeyspaceCreatingInt SchemaTestUtils.truncate(UserToken.class, cassandraTemplate); } + /** + * Post-process the {@link AsyncCassandraTemplate} before running the tests. + * + * @param template + */ + void prepareTemplate(AsyncCassandraTemplate template) { + template.setUsePreparedStatements(false); + } + @Test // DATACASS-343 void shouldSelectByQueryWithSorting() { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplatePreparedStatementsIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplatePreparedStatementsIntegrationTests.java new file mode 100644 index 000000000..a92aee267 --- /dev/null +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplatePreparedStatementsIntegrationTests.java @@ -0,0 +1,30 @@ +/* + * Copyright 2020 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 + * + * https://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.data.cassandra.core; + +/** + * Integration tests for {@link AsyncCassandraTemplate} with + * {@link AsyncCassandraTemplate#setUsePreparedStatements(boolean) prepared statements enabled}. + * + * @author Mark Paluch + */ +class AsyncCassandraTemplatePreparedStatementsIntegrationTests extends AsyncCassandraTemplateIntegrationTests { + + @Override + void prepareTemplate(AsyncCassandraTemplate template) { + template.setUsePreparedStatements(true); + } +} diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateUnitTests.java index 3138e3897..e62160dec 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateUnitTests.java @@ -87,6 +87,7 @@ public class AsyncCassandraTemplateUnitTests { void setUp() { template = new AsyncCassandraTemplate(session); + template.setUsePreparedStatements(false); when(session.executeAsync(any(Statement.class))).thenReturn(new TestResultSetFuture(resultSet)); when(row.getColumnDefinitions()).thenReturn(columnDefinitions); diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateIntegrationTests.java index 1bbf398e2..c7fe99d2c 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateIntegrationTests.java @@ -97,6 +97,8 @@ class CassandraTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrat template = new CassandraTemplate(new CqlTemplate(session), converter); + prepareTemplate(template); + SchemaTestUtils.potentiallyCreateTableFor(User.class, template); SchemaTestUtils.potentiallyCreateTableFor(UserToken.class, template); SchemaTestUtils.potentiallyCreateTableFor(BookReference.class, template); @@ -120,6 +122,15 @@ class CassandraTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrat SchemaTestUtils.truncate(WithMappedUdtList.class, template); } + /** + * Post-process the {@link CassandraTemplate} before running the tests. + * + * @param template + */ + void prepareTemplate(CassandraTemplate template) { + + } + @Test // DATACASS-343 void shouldSelectByQueryWithAllowFiltering() { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplatePreparedStatementIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplatePreparedStatementIntegrationTests.java new file mode 100644 index 000000000..8fecc9fcf --- /dev/null +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplatePreparedStatementIntegrationTests.java @@ -0,0 +1,30 @@ +/* + * Copyright 2020 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 + * + * https://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.data.cassandra.core; + +/** + * Integration tests for {@link CassandraTemplate} with {@link CassandraTemplate#setUsePreparedStatements(boolean) + * prepared statements enabled}. + * + * @author Mark Paluch + */ +class CassandraTemplatePreparedStatementIntegrationTests extends CassandraTemplateIntegrationTests { + + @Override + void prepareTemplate(CassandraTemplate template) { + template.setUsePreparedStatements(true); + } +} diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java index eb28c34c6..8384d435c 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java @@ -82,6 +82,7 @@ class CassandraTemplateUnitTests { void setUp() { template = new CassandraTemplate(session); + template.setUsePreparedStatements(false); when(session.execute(any(Statement.class))).thenReturn(resultSet); when(row.getColumnDefinitions()).thenReturn(columnDefinitions); diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateIntegrationTests.java index 3be5990ef..15008cb6f 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateIntegrationTests.java @@ -59,6 +59,7 @@ class ReactiveCassandraTemplateIntegrationTests extends AbstractKeyspaceCreating DefaultBridgedReactiveSession session = new DefaultBridgedReactiveSession(this.session); template = new ReactiveCassandraTemplate(new ReactiveCqlTemplate(session), converter); + prepareTemplate(template); SchemaTestUtils.potentiallyCreateTableFor(User.class, cassandraTemplate); SchemaTestUtils.potentiallyCreateTableFor(UserToken.class, cassandraTemplate); @@ -66,6 +67,15 @@ class ReactiveCassandraTemplateIntegrationTests extends AbstractKeyspaceCreating SchemaTestUtils.truncate(UserToken.class, cassandraTemplate); } + /** + * Post-process the {@link ReactiveCassandraTemplate} before running the tests. + * + * @param template + */ + void prepareTemplate(ReactiveCassandraTemplate template) { + template.setUsePreparedStatements(false); + } + @Test // DATACASS-335 void insertShouldInsertEntity() { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplatePreparedStatementIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplatePreparedStatementIntegrationTests.java new file mode 100644 index 000000000..0d8f18924 --- /dev/null +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplatePreparedStatementIntegrationTests.java @@ -0,0 +1,30 @@ +/* + * Copyright 2020 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 + * + * https://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.data.cassandra.core; + +/** + * Integration tests for {@link ReactiveCassandraTemplate} with + * {@link ReactiveCassandraTemplate#setUsePreparedStatements(boolean) prepared statements enabled}. + * + * @author Mark Paluch + */ +class ReactiveCassandraTemplatePreparedStatementIntegrationTests extends ReactiveCassandraTemplateIntegrationTests { + + @Override + void prepareTemplate(ReactiveCassandraTemplate template) { + template.setUsePreparedStatements(true); + } +} diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java index 6e259c457..b3397e4db 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java @@ -84,6 +84,7 @@ class ReactiveCassandraTemplateUnitTests { void setUp() { template = new ReactiveCassandraTemplate(session); + template.setUsePreparedStatements(false); when(session.execute(any(Statement.class))).thenReturn(Mono.just(reactiveResultSet)); when(row.getColumnDefinitions()).thenReturn(columnDefinitions); diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/example/CassandraTemplateExamples.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/example/CassandraTemplateExamples.java new file mode 100644 index 000000000..54b509492 --- /dev/null +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/example/CassandraTemplateExamples.java @@ -0,0 +1,49 @@ +/* + * Copyright 2020 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 + * + * https:://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.data.cassandra.example; + +import static org.springframework.data.cassandra.core.query.Criteria.*; +import static org.springframework.data.cassandra.core.query.Query.*; + +import org.springframework.data.cassandra.core.CassandraTemplate; + +import com.datastax.oss.driver.api.core.cql.SimpleStatement; + +/** + * @author Mark Paluch + */ +// @formatter:off +public class CassandraTemplateExamples { + + private CassandraTemplate template = null; + + void examples() { + // tag::preparedStatement[] + template.setUsePreparedStatements(true); + + Actor actorByQuery = template.selectOne(query(where("id").is(42)), Actor.class); + + Actor actorByStatement = template.selectOne( + SimpleStatement.newInstance("SELECT id, name FROM actor WHERE id = ?", 42), + Actor.class); + // end::preparedStatement[] + } + + static class Actor { + + } + +} diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/example/CqlTemplateExamples.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/example/CqlTemplateExamples.java index d569f40df..ea18bc923 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/example/CqlTemplateExamples.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/example/CqlTemplateExamples.java @@ -74,6 +74,13 @@ public class CqlTemplateExamples { } }); // end::listOfRowMapper[] + + // tag::preparedStatement[] + List lastNames = cqlTemplate.query( + session -> session.prepare("SELECT last_name FROM t_actor WHERE id = ?"), + ps -> ps.bind(1212L), + (row, rowNum) -> row.getString(0)); + // end::preparedStatement[] } // tag::findAllActors[] diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/isolated/RepositoryQueryMethodParameterTypesIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/isolated/RepositoryQueryMethodParameterTypesIntegrationTests.java index 77e1d71d5..082875564 100755 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/isolated/RepositoryQueryMethodParameterTypesIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/isolated/RepositoryQueryMethodParameterTypesIntegrationTests.java @@ -36,6 +36,7 @@ import org.springframework.context.annotation.Configuration; import org.springframework.core.convert.converter.Converter; import org.springframework.data.cassandra.CassandraInvalidQueryException; import org.springframework.data.cassandra.config.SchemaAction; +import org.springframework.data.cassandra.core.CassandraAdminTemplate; import org.springframework.data.cassandra.core.convert.CassandraCustomConversions; import org.springframework.data.cassandra.core.convert.MappingCassandraConverter; import org.springframework.data.cassandra.core.mapping.CassandraMappingContext; @@ -75,6 +76,13 @@ class RepositoryQueryMethodParameterTypesIntegrationTests public SchemaAction getSchemaAction() { return SchemaAction.RECREATE_DROP_UNUSED; } + + @Override + public CassandraAdminTemplate cassandraTemplate() { + CassandraAdminTemplate template = super.cassandraTemplate(); + template.setUsePreparedStatements(false); + return template; + } } @Autowired AllPossibleTypesRepository allPossibleTypesRepository; diff --git a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/repository/CoroutineRepositoryUnitTests.kt b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/repository/CoroutineRepositoryUnitTests.kt index e97a7c728..9de48d85e 100644 --- a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/repository/CoroutineRepositoryUnitTests.kt +++ b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/repository/CoroutineRepositoryUnitTests.kt @@ -15,7 +15,6 @@ */ package org.springframework.data.cassandra.repository -import com.datastax.oss.driver.api.core.cql.Statement import io.mockk.every import io.mockk.mockk import kotlinx.coroutines.runBlocking @@ -55,7 +54,7 @@ class CoroutineRepositoryUnitTests { fun `should discard result of suspended query method without result`() { every { resultSet.wasApplied() } returns true - every { cqlOperations.queryForResultSet(any>()) } returns Mono.just(resultSet) + every { operations.execute(any()) } returns Mono.just(resultSet) val repository = repositoryFactory.getRepository(PersonRepository::class.java) diff --git a/src/main/asciidoc/new-features.adoc b/src/main/asciidoc/new-features.adoc index 4cf175db0..81bf54bfa 100644 --- a/src/main/asciidoc/new-features.adoc +++ b/src/main/asciidoc/new-features.adoc @@ -3,6 +3,11 @@ This chapter summarizes changes and new features for each release. +[[new-features.3-2-0]] +== What's new in Spring Data for Apache Cassandra 3.2 + +* <> using `CassandraTemplate` and repositories (enabled by default). + [[new-features.3-1-0]] == What's new in Spring Data for Apache Cassandra 3.1 diff --git a/src/main/asciidoc/reference/cassandra-repositories.adoc b/src/main/asciidoc/reference/cassandra-repositories.adoc index bb63ec84d..149c7ddbc 100644 --- a/src/main/asciidoc/reference/cassandra-repositories.adoc +++ b/src/main/asciidoc/reference/cassandra-repositories.adoc @@ -3,6 +3,7 @@ This chapter covers the details of the Spring Data Repository support for Apache Cassandra. Cassandra's repository support builds on the core repository support explained in "`<>`". +Cassandra repositories use `CassandraTemplate` and its wired `CqlTemplate` as infrastructure beans. You should understand the basic concepts explained there before proceeding. [[cassandra-repo-usage]] diff --git a/src/main/asciidoc/reference/cassandra.adoc b/src/main/asciidoc/reference/cassandra.adoc index 74601b5e4..ac72c7f81 100644 --- a/src/main/asciidoc/reference/cassandra.adoc +++ b/src/main/asciidoc/reference/cassandra.adoc @@ -520,7 +520,7 @@ The `CqlTemplate` class executes CQL queries and update statements, performs ite It also catches CQL exceptions and translates them to the generic, more informative, exception hierarchy defined in the `org.springframework.dao` package. When you use the `CqlTemplate` for your code, you need only implement callback interfaces, which have a clearly defined contract. -Given a `Connection`, the `PreparedStatementCreator` callback interface creates a prepared statement with the provided CQL and any necessary parameter arguments. +Given a `Connection`, the `PreparedStatementCreator` callback interface creates a <> with the provided CQL and any necessary parameter arguments. The `RowCallbackHandler` interface extracts values from each row of a `ResultSet`. The `CqlTemplate` can be used within a DAO implementation through direct instantiation with a `SessionFactory` reference or be configured in the Spring container and given to DAOs as a bean reference. `CqlTemplate` is a foundational building block for <>. @@ -1118,3 +1118,78 @@ The terminating methods (`first()`, `one()`, `all()`, and `stream()`) handle swi WARNING: The new fluent template API methods (that is, `query(..)`, `insert(..)`, `update(..)`, and `delete(..)`) use effectively thread-safe supporting objects to compose the CQL statement. However, it comes at the added cost of additional young-gen JVM heap overhead, since the design is based on final fields for the various CQL statement components and construction on mutation. You should be careful when possibly inserting or deleting a large number of objects (such as inside of a loop, for instance). + +[[cassandra.template.prepared-statements]] +== Prepared Statements + +CQL statements that are executed multiple times can be prepared and stored in a `PreparedStatement` object to improve query performance. +Both, the driver and Cassandra maintain a mapping of `PreparedStatement` queries to their metadata. +You can use prepared statements through the following abstractions: + +* `CqlTemplate` through the choice of API +* `CassandraTemplate` by enabling prepared statements +* Cassandra repositories as they are built on `CassandraTemplate` + +[[cassandra.template.prepared-statements.cql]] +=== Using `CqlTemplate` + +The `CqlTemplate` class (and its asynchronous and reactive variants) offers various methods accepting static CQL, `Statement` objects and `PreparedStatementCreator`. +Methods accepting static CQL without additional arguments typically run the CQL statement as-is without further processing. +Methods accepting static CQL in combination with an arguments array (such as `execute(String cql, Object... args)` and `queryForRows(String cql, Object... args)`) use prepared statements. +Internally, these methods create a `PreparedStatementCreator` and `PreparedStatementBinder` objects to prepare the statement and later on to bind values to the statement to run it. +Spring Data Cassandra generally uses index-based parameter bindings for prepared statements. + +Since Cassandra Driver version 4, prepared statements are cached on the driver level which removes the need to keep track of prepared statements in the application. + +The following example shows how to issue a query with a parametrized prepared statement: + +==== +[source,java,indent=0] +---- +include::../{example-root}/CqlTemplateExamples.java[tags=lastName] +---- +==== + +In cases where you require more control over statement preparation and parameter binding (for example, using named binding parameters), you can fully control prepared statement creation and parameter binding by calling query methods with `PreparedStatementCreator` and `PreparedStatementBinder` arguments: + +==== +[source,java,indent=0] +---- +include::../{example-root}/CqlTemplateExamples.java[tags=preparedStatement] +---- +==== + +Spring Data Cassandra ships with classes supporting that pattern in the `cql` package: + +* `SimplePreparedStatementCreator` - utility class to create a prepared statement. +* `ArgumentPreparedStatementBinder` - utility class to bind arguments to a prepared statement. + +[[cassandra.template.prepared-statements.cassandra-template]] +=== Using `CassandraTemplate` + +The `CassandraTemplate` class is built on top of `CqlTemplate` to provide a higher level of abstraction. +The use of prepared statements can be controlled directly on `CassandraTemplate` (and its asynchronous and reactive variants) by calling `setUsePreparedStatements(false)` respective `setUsePreparedStatements(true)`. +Note that the use of prepared statements by `CassandraTemplate` is enabled by default. + +The following example shows the use of methods that generate and that accept CQL: + +==== +[source,java,indent=0] +---- +include::../{example-root}/CassandraTemplateExamples.java[tags=preparedStatement] +---- +==== + +Calling entity-bound methods such as `select(Query, Class)` or `update(Query, Update, Class)` build CQL statements themselves to perform the intended operations. +Some `CassandraTemplate` methods (such as `select(Statement, Class)`) also accepts CQL `Statement` objects as part of their API. + +It's possible to participate in prepared statements when calling methods accepting a `Statement` with a `SimpleStatement` object. +The template API extracts the query string and parameters (positional and named parameters) and uses these to prepare, bind, and run the statement. +Non-``SimpleStatement`` objects cannot be used with prepared statements. + +[[cassandra.template.prepared-statements.caching]] +=== Caching Prepared Statements + +Since Cassandra driver 4.0, prepared statements are cached by the `CqlSession` cache so it is okay to prepare the same string twice. +Previous versions required caching of prepared statements outside of the driver. +See also the https://docs.datastax.com/en/developer/java-driver/latest/manual/core/statements/prepared/[Driver documentation on Prepared Statements] for further reference. diff --git a/src/main/asciidoc/reference/reactive-cassandra-repositories.adoc b/src/main/asciidoc/reference/reactive-cassandra-repositories.adoc index 45e88027e..15b71ad4a 100644 --- a/src/main/asciidoc/reference/reactive-cassandra-repositories.adoc +++ b/src/main/asciidoc/reference/reactive-cassandra-repositories.adoc @@ -4,6 +4,8 @@ This chapter outlines the specialties handled by the reactive repository support for Apache Cassandra. It builds on the core repository infrastructure explained in <>, so you should have a good understanding of the basic concepts explained there. +Cassandra repositories use `ReactiveCassandraTemplate` and its wired `ReactiveCqlTemplate` as infrastructure beans. + Reactive usage is broken up into two phases: Composition and Execution. Calling repository methods lets you compose a reactive sequence by obtaining `Publisher` instances and applying operators. diff --git a/src/main/asciidoc/reference/reactive-cassandra.adoc b/src/main/asciidoc/reference/reactive-cassandra.adoc index 565280dc5..47ffc2adf 100644 --- a/src/main/asciidoc/reference/reactive-cassandra.adoc +++ b/src/main/asciidoc/reference/reactive-cassandra.adoc @@ -150,7 +150,7 @@ The `ReactiveCqlTemplate` class runs CQL queries and update statements and perfo It also catches CQL exceptions and translates them into the generic, more informative, exception hierarchy defined in the `org.springframework.dao` package. When you use the `ReactiveCqlTemplate` in your code, you need only implement callback interfaces, which have a clearly defined contract. -Given a `Connection`, the `ReactivePreparedStatementCreator` callback interface creates a prepared statement with the provided CQL and any necessary parameter arguments. +Given a `Connection`, the `ReactivePreparedStatementCreator` callback interface creates a <> with the provided CQL and any necessary parameter arguments. The `RowCallbackHandler` interface extracts values from each row of a `ReactiveResultSet`.