From b79f8a4caabbd3db2d28971064a84f1d0ef04d94 Mon Sep 17 00:00:00 2001 From: tomekl007 Date: Wed, 5 Aug 2020 12:36:05 +0200 Subject: [PATCH] DATACASS-767 - Allow configuring keyspace per Statement. Original pull request: #177. --- .../data/cassandra/core/DeleteOptions.java | 12 ++++ .../data/cassandra/core/InsertOptions.java | 12 ++++ .../data/cassandra/core/UpdateOptions.java | 12 ++++ .../cassandra/core/cql/CassandraAccessor.java | 52 +++++++++++++-- .../data/cassandra/core/cql/CqlTemplate.java | 12 ++-- .../data/cassandra/core/cql/QueryOptions.java | 46 ++++++++++++- .../cassandra/core/cql/QueryOptionsUtil.java | 20 +++++- .../core/cql/ReactiveCqlTemplate.java | 64 +++++++++++++++---- .../data/cassandra/core/cql/WriteOptions.java | 16 ++++- .../CassandraTemplateIntegrationTests.java | 32 ++++++++++ .../core/cql/CassandraAccessorUnitTests.java | 10 +++ .../core/cql/CqlTemplateIntegrationTests.java | 33 +++++++++- .../core/cql/CqlTemplateUnitTests.java | 30 ++++++++- .../core/cql/QueryOptionsUnitTests.java | 8 ++- .../core/cql/QueryOptionsUtilUnitTests.java | 43 +++++++++++++ .../ReactiveCqlTemplateIntegrationTests.java | 37 ++++++++++- .../cql/ReactiveCqlTemplateUnitTests.java | 19 +++++- .../test/util/CassandraDelegate.java | 20 +++++- 18 files changed, 440 insertions(+), 38 deletions(-) diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/DeleteOptions.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/DeleteOptions.java index 180ef7248..bb8f61e28 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/DeleteOptions.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/DeleteOptions.java @@ -28,11 +28,13 @@ import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; import com.datastax.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.CqlIdentifier; /** * Extension to {@link WriteOptions} for use with {@code DELETE} operations. * * @author Mark Paluch + * @author Tomasz Lelek * @since 2.2 */ public class DeleteOptions extends WriteOptions { @@ -305,6 +307,16 @@ public class DeleteOptions extends WriteOptions { return this; } + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.cql.WriteOptions.WriteOptionsBuilder#keyspace() + */ + @Override + public DeleteOptionsBuilder keyspace(CqlIdentifier keyspace) { + + super.keyspace(keyspace); + return this; + } + /** * Use light-weight transactions by applying {@code IF EXISTS}. Replaces a previous {@link #ifCondition(Filter)}. * diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/InsertOptions.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/InsertOptions.java index b432a43ce..648a15cad 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/InsertOptions.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/InsertOptions.java @@ -24,12 +24,14 @@ import org.springframework.data.cassandra.core.cql.WriteOptions; import org.springframework.lang.Nullable; import com.datastax.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.CqlIdentifier; /** * Extension to {@link WriteOptions} for use with {@code INSERT} operations. * * @author Mark Paluch * @author Lukasz Antoniak + * @author Tomasz Lelek * @since 2.0 */ public class InsertOptions extends WriteOptions { @@ -304,6 +306,16 @@ public class InsertOptions extends WriteOptions { return this; } + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.cql.WriteOptions.WriteOptionsBuilder#keyspace() + */ + @Override + public InsertOptionsBuilder keyspace(CqlIdentifier keyspace) { + + super.keyspace(keyspace); + return this; + } + /** * Use light-weight transactions by applying {@code IF NOT EXISTS}. * diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/UpdateOptions.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/UpdateOptions.java index bf8e13ab4..7b7a4691c 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/UpdateOptions.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/UpdateOptions.java @@ -28,12 +28,14 @@ import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; import com.datastax.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.CqlIdentifier; /** * Extension to {@link WriteOptions} for use with {@code UPDATE} operations. * * @author Mark Paluch * @author Lukasz Antoniak + * @author Tomasz Lelek * @since 2.0 */ public class UpdateOptions extends WriteOptions { @@ -312,6 +314,16 @@ public class UpdateOptions extends WriteOptions { return this; } + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.cql.WriteOptions.WriteOptionsBuilder#keyspace() + */ + @Override + public UpdateOptionsBuilder keyspace(CqlIdentifier keyspace) { + + super.keyspace(keyspace); + return this; + } + /** * Use light-weight transactions by applying {@code IF EXISTS}. Replaces a previous {@link #ifCondition(Filter)}. * diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraAccessor.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraAccessor.java index 08513b1f7..2391fe46d 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraAccessor.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraAccessor.java @@ -18,11 +18,6 @@ package org.springframework.data.cassandra.core.cql; import java.util.Map; import java.util.Optional; -import com.datastax.oss.driver.api.core.ConsistencyLevel; -import com.datastax.oss.driver.api.core.CqlSession; -import com.datastax.oss.driver.api.core.cql.SimpleStatement; -import com.datastax.oss.driver.api.core.cql.Statement; -import com.datastax.oss.driver.api.core.retry.RetryPolicy; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -33,6 +28,15 @@ import org.springframework.data.cassandra.core.cql.session.DefaultSessionFactory import org.springframework.lang.Nullable; import org.springframework.util.Assert; +import com.datastax.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.CqlIdentifier; +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.cql.BatchStatement; +import com.datastax.oss.driver.api.core.cql.BoundStatement; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.api.core.retry.RetryPolicy; + /** * {@link CassandraAccessor} provides access to a Cassandra {@link SessionFactory} and the * {@link CassandraExceptionTranslator}. @@ -43,6 +47,7 @@ import org.springframework.util.Assert; * @author David Webb * @author Mark Paluch * @author John Blum + * @author Tomasz Lelek * @see org.springframework.beans.factory.InitializingBean * @see com.datastax.oss.driver.api.core.CqlSession */ @@ -79,6 +84,12 @@ public class CassandraAccessor implements InitializingBean { private @Nullable SessionFactory sessionFactory; + /** + * If this variable is set to a value, it will be used for setting the {@code keyspace} property on statements used + * for query processing. + */ + private @Nullable CqlIdentifier keyspace; + /** * Ensures the Cassandra {@link CqlSession} and exception translator has been propertly set. */ @@ -286,6 +297,24 @@ public class CassandraAccessor implements InitializingBean { return this.sessionFactory; } + /** + * Set the keyspace for this template. If it is null, then the default {@link CqlSession} level keyspace will be used. + * + * @see SimpleStatement#setKeyspace(CqlIdentifier) + * @see BatchStatement#setKeyspace(CqlIdentifier) + */ + public void setKeyspace(@Nullable CqlIdentifier keyspace) { + this.keyspace = keyspace; + } + + /** + * @return the {@link CqlIdentifier} keyspace for this template. + */ + @Nullable + public CqlIdentifier getKeyspace() { + return this.keyspace; + } + /** * Create a {@link SimpleStatement} given {@code cql}. * @@ -312,6 +341,7 @@ public class CassandraAccessor implements InitializingBean { Statement statementToUse = statement; ConsistencyLevel consistencyLevel = getConsistencyLevel(); ConsistencyLevel serialConsistencyLevel = getSerialConsistencyLevel(); + CqlIdentifier keyspace = getKeyspace(); int pageSize = getPageSize(); if (consistencyLevel != null) { @@ -326,6 +356,18 @@ public class CassandraAccessor implements InitializingBean { statementToUse = statementToUse.setPageSize(pageSize); } + if (keyspace != null) { + if (statementToUse instanceof BoundStatement) { + throw new IllegalArgumentException("Keyspace cannot be set for a BoundStatement"); + } + if (statementToUse instanceof BatchStatement) { + statementToUse = ((BatchStatement) statementToUse).setKeyspace(keyspace); + } + if (statementToUse instanceof SimpleStatement) { + statementToUse = ((SimpleStatement) statementToUse).setKeyspace(keyspace); + } + } + statementToUse = getExecutionProfileResolver().apply(statementToUse); return statementToUse; diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CqlTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CqlTemplate.java index 01a3c71dd..227f29652 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CqlTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CqlTemplate.java @@ -21,6 +21,12 @@ import java.util.Map; import java.util.function.Function; import java.util.stream.StreamSupport; +import org.springframework.dao.DataAccessException; +import org.springframework.dao.support.DataAccessUtils; +import org.springframework.data.cassandra.SessionFactory; +import org.springframework.lang.Nullable; +import org.springframework.util.Assert; + 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; @@ -29,12 +35,6 @@ import com.datastax.oss.driver.api.core.cql.Row; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.core.metadata.Node; -import org.springframework.dao.DataAccessException; -import org.springframework.dao.support.DataAccessUtils; -import org.springframework.data.cassandra.SessionFactory; -import org.springframework.lang.Nullable; -import org.springframework.util.Assert; - /** * This is the central class in the CQL core package. It simplifies the use of CQL and helps to avoid common * errors. It executes core CQL workflow, leaving application code to provide CQL and extract results. This class diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptions.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptions.java index 602b4d669..b16a13adc 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptions.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptions.java @@ -23,6 +23,8 @@ import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; import com.datastax.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.CqlIdentifier; +import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.config.DriverExecutionProfile; /** @@ -31,6 +33,7 @@ import com.datastax.oss.driver.api.core.config.DriverExecutionProfile; * * @author David Webb * @author Mark Paluch + * @author Tomasz Lelek */ public class QueryOptions { @@ -48,9 +51,17 @@ public class QueryOptions { private final @Nullable Boolean tracing; + private final @Nullable CqlIdentifier keyspace; + protected QueryOptions(@Nullable ConsistencyLevel consistencyLevel, ExecutionProfileResolver executionProfileResolver, @Nullable Integer pageSize, @Nullable ConsistencyLevel serialConsistencyLevel, Duration timeout, @Nullable Boolean tracing) { + this(consistencyLevel, executionProfileResolver, pageSize, serialConsistencyLevel, timeout, tracing, null); + } + + protected QueryOptions(@Nullable ConsistencyLevel consistencyLevel, ExecutionProfileResolver executionProfileResolver, + @Nullable Integer pageSize, @Nullable ConsistencyLevel serialConsistencyLevel, Duration timeout, + @Nullable Boolean tracing, @Nullable CqlIdentifier keyspace) { this.consistencyLevel = consistencyLevel; this.executionProfileResolver = executionProfileResolver; @@ -58,6 +69,7 @@ public class QueryOptions { this.serialConsistencyLevel = serialConsistencyLevel; this.timeout = timeout; this.tracing = tracing; + this.keyspace = keyspace; } /** @@ -154,6 +166,15 @@ public class QueryOptions { return this.tracing; } + /** + * @return the keyspace associated with the query. If it is null, it means that the default keyspace from + * {@link CqlSession} will be used. + */ + @Nullable + public CqlIdentifier getKeyspace() { + return keyspace; + } + /* * (non-Javadoc) * @see java.lang.Object#equals(java.lang.Object) @@ -191,7 +212,11 @@ public class QueryOptions { return false; } - return ObjectUtils.nullSafeEquals(tracing, options.tracing); + if (!ObjectUtils.nullSafeEquals(tracing, options.tracing)) { + return false; + } + + return ObjectUtils.nullSafeEquals(keyspace, options.keyspace); } /* @@ -206,6 +231,7 @@ public class QueryOptions { result = 31 * result + ObjectUtils.nullSafeHashCode(serialConsistencyLevel); result = 31 * result + ObjectUtils.nullSafeHashCode(timeout); result = 31 * result + ObjectUtils.nullSafeHashCode(tracing); + result = 31 * result + ObjectUtils.nullSafeHashCode(keyspace); return result; } @@ -229,6 +255,8 @@ public class QueryOptions { protected @Nullable Boolean tracing; + protected @Nullable CqlIdentifier keyspace; + QueryOptionsBuilder() {} QueryOptionsBuilder(QueryOptions queryOptions) { @@ -239,6 +267,7 @@ public class QueryOptions { this.serialConsistencyLevel = queryOptions.serialConsistencyLevel; this.timeout = queryOptions.timeout; this.tracing = queryOptions.tracing; + this.keyspace = queryOptions.keyspace; } /** @@ -430,6 +459,19 @@ public class QueryOptions { return tracing(true); } + /** + * Sets the keyspace to use. + * + * @param keyspace if it is null, then the default {@link CqlSession} level keyspace will be used. + * @return {@code this} {@link QueryOptionsBuilder} + */ + public QueryOptionsBuilder keyspace(CqlIdentifier keyspace) { + + this.keyspace = keyspace; + + return this; + } + /** * Builds a new {@link QueryOptions} with the configured values. * @@ -437,7 +479,7 @@ public class QueryOptions { */ public QueryOptions build() { return new QueryOptions(this.consistencyLevel, this.executionProfileResolver, this.pageSize, - this.serialConsistencyLevel, this.timeout, this.tracing); + this.serialConsistencyLevel, this.timeout, this.tracing, this.keyspace); } } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtil.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtil.java index 3fac279aa..093d4b5ed 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtil.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtil.java @@ -17,6 +17,11 @@ package org.springframework.data.cassandra.core.cql; import java.time.Duration; +import org.springframework.util.Assert; + +import com.datastax.oss.driver.api.core.cql.BatchStatement; +import com.datastax.oss.driver.api.core.cql.BoundStatement; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.querybuilder.delete.Delete; import com.datastax.oss.driver.api.querybuilder.delete.DeleteSelection; @@ -24,13 +29,12 @@ import com.datastax.oss.driver.api.querybuilder.insert.Insert; import com.datastax.oss.driver.api.querybuilder.update.Update; import com.datastax.oss.driver.api.querybuilder.update.UpdateStart; -import org.springframework.util.Assert; - /** * Utility class to associate {@link QueryOptions} and {@link WriteOptions} with QueryBuilder {@link Statement}s. * * @author Mark Paluch * @author Lukasz Antoniak + * @author Tomasz Lelek * @since 2.0 */ public abstract class QueryOptionsUtil { @@ -73,6 +77,18 @@ public abstract class QueryOptionsUtil { // statement wrapped in the conditional null check to avoid additional garbage and added GC pressure. statementToUse = statementToUse.setTracing(Boolean.TRUE.equals(queryOptions.getTracing())); } + if (queryOptions.getKeyspace() != null) { + if (statementToUse instanceof BoundStatement) { + throw new IllegalArgumentException("Keyspace cannot be set for a BoundStatement"); + } + if (statementToUse instanceof BatchStatement) { + statementToUse = ((BatchStatement) statementToUse).setKeyspace(queryOptions.getKeyspace()); + } + if (statementToUse instanceof SimpleStatement) { + statementToUse = ((SimpleStatement) statementToUse).setKeyspace(queryOptions.getKeyspace()); + } + + } return (T) statementToUse; } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplate.java index c9d5d437d..2e185833b 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplate.java @@ -15,22 +15,14 @@ */ package org.springframework.data.cassandra.core.cql; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + import java.util.Map; import java.util.Optional; import java.util.function.Function; -import com.datastax.oss.driver.api.core.ConsistencyLevel; -import com.datastax.oss.driver.api.core.DriverException; -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.Row; -import com.datastax.oss.driver.api.core.cql.SimpleStatement; -import com.datastax.oss.driver.api.core.cql.Statement; -import com.datastax.oss.driver.api.core.retry.RetryPolicy; import org.reactivestreams.Publisher; -import reactor.core.publisher.Flux; -import reactor.core.publisher.Mono; - import org.springframework.dao.DataAccessException; import org.springframework.dao.support.DataAccessUtils; import org.springframework.data.cassandra.ReactiveResultSet; @@ -40,6 +32,18 @@ import org.springframework.data.cassandra.core.cql.session.DefaultReactiveSessio import org.springframework.lang.Nullable; import org.springframework.util.Assert; +import com.datastax.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.CqlIdentifier; +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.DriverException; +import com.datastax.oss.driver.api.core.cql.BatchStatement; +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.Row; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.api.core.retry.RetryPolicy; + /** * This is the central class in the CQL core package for reactive Cassandra data access. It simplifies the use of * CQL and helps to avoid common errors. It executes core CQL workflow, leaving application code to provide CQL and @@ -68,6 +72,7 @@ import org.springframework.util.Assert; * NOTE: An instance of this class is thread-safe once configured. * * @author Mark Paluch + * @author Tomasz Lelek * @since 2.0 * @see PreparedStatementCreator * @see PreparedStatementBinder @@ -103,6 +108,12 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re */ private @Nullable ConsistencyLevel serialConsistencyLevel; + /** + * If this variable is set to a value, it will be used for setting the {@code keyspace} property on statements used + * for query processing. + */ + private @Nullable CqlIdentifier keyspace; + /** * Construct a new {@link ReactiveCqlTemplate Note: The {@link ReactiveSessionFactory} has to be set before using the * instance. @@ -254,6 +265,24 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return this.serialConsistencyLevel; } + /** + * Set the keyspace for this template. If it is null, then the default {@link CqlSession} level keyspace will be used. + * + * @see SimpleStatement#setKeyspace(CqlIdentifier) + * @see BatchStatement#setKeyspace(CqlIdentifier) + */ + public void setKeyspace(@Nullable CqlIdentifier keyspace) { + this.keyspace = keyspace; + } + + /** + * @return the {@link CqlIdentifier} keyspace for this template. + */ + @Nullable + public CqlIdentifier getKeyspace() { + return this.keyspace; + } + // ------------------------------------------------------------------------- // Methods dealing with a plain org.springframework.data.cassandra.core.cql.ReactiveSession // ------------------------------------------------------------------------- @@ -820,6 +849,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re Statement statementToUse = statement; ConsistencyLevel consistencyLevel = getConsistencyLevel(); ConsistencyLevel serialConsistencyLevel = getSerialConsistencyLevel(); + CqlIdentifier keyspace = getKeyspace(); int pageSize = getPageSize(); if (consistencyLevel != null) { @@ -834,6 +864,18 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re statementToUse = statementToUse.setPageSize(pageSize); } + if (keyspace != null) { + if (statementToUse instanceof BoundStatement) { + throw new IllegalArgumentException("Keyspace cannot be set for a BoundStatement"); + } + if (statementToUse instanceof BatchStatement) { + statementToUse = ((BatchStatement) statementToUse).setKeyspace(keyspace); + } + if (statementToUse instanceof SimpleStatement) { + statementToUse = ((SimpleStatement) statementToUse).setKeyspace(keyspace); + } + } + statementToUse = getExecutionProfileResolver().apply(statementToUse); return statementToUse; diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/WriteOptions.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/WriteOptions.java index 365005805..a6ee1baeb 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/WriteOptions.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/WriteOptions.java @@ -19,12 +19,13 @@ import java.time.Duration; import java.time.Instant; import java.util.concurrent.TimeUnit; -import com.datastax.oss.driver.api.core.ConsistencyLevel; - import org.springframework.lang.Nullable; import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; +import com.datastax.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.CqlIdentifier; + /** * Cassandra Write Options are an extension to {@link QueryOptions} for write operations. {@link WriteOptions} allow * tuning of various query options on a per-request level. Only options that are set are applied to queries. @@ -32,6 +33,7 @@ import org.springframework.util.ObjectUtils; * @author David Webb * @author Mark Paluch * @author Lukasz Antoniak + * @author Tomasz Lelek * @see QueryOptions */ public class WriteOptions extends QueryOptions { @@ -270,6 +272,16 @@ public class WriteOptions extends QueryOptions { return this; } + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.cql.QueryOptions.QueryOptionsBuilder#keyspace() + */ + @Override + public WriteOptionsBuilder keyspace(CqlIdentifier keyspace) { + + super.keyspace(keyspace); + return this; + } + /** * Sets the time to live in seconds for write operations. * 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 670913eeb..8dde813de 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 @@ -38,9 +38,11 @@ import java.util.stream.Stream; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.springframework.data.annotation.Id; +import org.springframework.data.cassandra.CassandraInvalidQueryException; import org.springframework.data.cassandra.core.convert.MappingCassandraConverter; import org.springframework.data.cassandra.core.cql.CqlTemplate; import org.springframework.data.cassandra.core.cql.PrimaryKeyType; +import org.springframework.data.cassandra.core.cql.QueryOptions; import org.springframework.data.cassandra.core.mapping.BasicMapId; import org.springframework.data.cassandra.core.mapping.Embedded; import org.springframework.data.cassandra.core.mapping.PrimaryKey; @@ -70,10 +72,12 @@ import com.datastax.oss.driver.api.core.uuid.Uuids; * * @author Mark Paluch * @author Christoph Strobl + * @author Tomasz Lelek */ class CassandraTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrationTests { private static final Version CASSANDRA_3 = Version.parse("3.0"); + private static final Version CASSANDRA_4 = Version.parse("4.0"); private Version cassandraVersion; @@ -432,6 +436,34 @@ class CassandraTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrat assertThat(template.delete(query, User.class)).isTrue(); } + @Test // DATACASS-767 + void selectByQueryWithKeyspaceShouldRetrieveData() { + assumeTrue(cassandraVersion.isGreaterThanOrEqualTo(CASSANDRA_4)); + + QueryOptions queryOptions = QueryOptions.builder().keyspace(CqlIdentifier.fromCql(keyspace)).build(); + + User user = new User("heisenberg", "Walter", "White"); + template.insert(user); + + Query query = Query.query(where("id").is("heisenberg")).queryOptions(queryOptions); + assertThat(template.select(query, User.class)).isNotEmpty(); + } + + @Test // DATACASS-767 + void selectByQueryWithNonExistingKeyspaceShouldThrowThatKeyspaceDoesNotExists() { + assumeTrue(cassandraVersion.isGreaterThanOrEqualTo(CASSANDRA_4)); + + QueryOptions queryOptions = QueryOptions.builder().keyspace(CqlIdentifier.fromCql("non_existing")).build(); + + User user = new User("heisenberg", "Walter", "White"); + template.insert(user); + + Query query = Query.query(where("id").is("heisenberg")).queryOptions(queryOptions); + assertThatThrownBy(() -> assertThat(template.select(query, User.class)).isEmpty()) + .isInstanceOf(CassandraInvalidQueryException.class) + .hasMessageContaining("Keyspace 'non_existing' does not exist"); + } + @Test // DATACASS-182 void stream() { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CassandraAccessorUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CassandraAccessorUnitTests.java index 75a5a4657..ff44e5b26 100755 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CassandraAccessorUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CassandraAccessorUnitTests.java @@ -20,9 +20,11 @@ import static org.assertj.core.api.Assertions.*; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; + import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.CqlSession; /** @@ -102,4 +104,12 @@ class CassandraAccessorUnitTests { assertThat(e).hasMessageContaining("SessionFactory was not properly initialized"); } } + + @Test // DATACASS-767 + void setAndGetKeyspace() { + CqlIdentifier keyspace = CqlIdentifier.fromCql("ks1"); + cassandraAccessor.setKeyspace(keyspace); + assertThat(cassandraAccessor.getKeyspace()).isEqualTo(keyspace); + } + } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CqlTemplateIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CqlTemplateIntegrationTests.java index ec590efe0..55f964fa4 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CqlTemplateIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CqlTemplateIntegrationTests.java @@ -16,6 +16,7 @@ package org.springframework.data.cassandra.core.cql; import static org.assertj.core.api.Assertions.*; +import static org.junit.Assume.*; import java.util.ArrayList; import java.util.List; @@ -25,19 +26,25 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.springframework.data.cassandra.CassandraInvalidQueryException; +import org.springframework.data.cassandra.support.CassandraVersion; import org.springframework.data.cassandra.test.util.AbstractKeyspaceCreatingIntegrationTests; +import org.springframework.data.util.Version; +import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.cql.SimpleStatement; /** * Integration tests for {@link CqlTemplate}. * * @author Mark Paluch + * @author Tomasz Lelek */ class CqlTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrationTests { - + private static final Version CASSANDRA_4 = Version.parse("4.0"); private static final AtomicBoolean initialized = new AtomicBoolean(); private CqlTemplate template; + private Version cassandraVersion; @BeforeEach void before() { @@ -51,6 +58,7 @@ class CqlTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrationTes template = new CqlTemplate(); template.setSession(getSession()); + cassandraVersion = CassandraVersion.get(session); } @Test // DATACASS-292 @@ -168,4 +176,27 @@ class CqlTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrationTes assertThat(map).containsEntry("id", "WHITE").containsEntry("username", "Walter"); } + + @Test // DATACASS-767 + void selectByQueryWithKeyspaceShouldRetrieveData() { + assumeTrue(cassandraVersion.isGreaterThanOrEqualTo(CASSANDRA_4)); + + template.setKeyspace(CqlIdentifier.fromCql(keyspace)); + + String id = template.queryForObject("SELECT id FROM user;", String.class); + + assertThat(id).isEqualTo("WHITE"); + } + + @Test // DATACASS-767 + void selectByQueryWithNonExistingKeyspaceShouldThrowThatKeyspaceDoesNotExists() { + assumeTrue(cassandraVersion.isGreaterThanOrEqualTo(CASSANDRA_4)); + + template.setKeyspace(CqlIdentifier.fromCql("non_existing")); + + assertThatThrownBy(() -> template.queryForObject("SELECT id FROM user;", String.class)) + .isInstanceOf(CassandraInvalidQueryException.class) + .hasMessageContaining("Keyspace 'non_existing' does not exist"); + + } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CqlTemplateUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CqlTemplateUnitTests.java index d230c53b0..9d07a882b 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CqlTemplateUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/CqlTemplateUnitTests.java @@ -28,12 +28,12 @@ import java.util.function.Consumer; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; + import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.junit.jupiter.MockitoSettings; import org.mockito.quality.Strictness; - import org.springframework.dao.EmptyResultDataAccessException; import org.springframework.dao.IncorrectResultSizeDataAccessException; import org.springframework.data.cassandra.CassandraConnectionFailureException; @@ -41,6 +41,7 @@ import org.springframework.data.cassandra.CassandraInvalidQueryException; import org.springframework.lang.Nullable; import com.datastax.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.DefaultConsistencyLevel; import com.datastax.oss.driver.api.core.NoNodeAvailableException; @@ -57,6 +58,7 @@ import com.datastax.oss.driver.api.core.servererrors.InvalidQueryException; * Unit tests for {@link CqlTemplate}. * * @author Mark Paluch + * @author Tomasz Lelek */ @ExtendWith(MockitoExtension.class) @MockitoSettings(strictness = Strictness.LENIENT) @@ -682,13 +684,30 @@ class CqlTemplateUnitTests { assertThat(applied).isTrue(); } + @Test // DATACASS-767 + void executeCqlWithKeyspaceShouldCallExecution() { + + doTestStrings(5, DefaultConsistencyLevel.ONE, null, "foo", cqlTemplate -> { + + cqlTemplate.execute("SELECT * from USERS"); + + verify(session).execute(any(Statement.class)); + }, CqlIdentifier.fromCql("some_keyspace")); + } + private void doTestStrings(Consumer cqlTemplateConsumer) { - doTestStrings(null, null, null, null, cqlTemplateConsumer); + doTestStrings(null, null, null, null, cqlTemplateConsumer, null); } private void doTestStrings(@Nullable Integer fetchSize, @Nullable ConsistencyLevel consistencyLevel, @Nullable ConsistencyLevel serialConsistencyLevel, @Nullable String executionProfile, Consumer cqlTemplateConsumer) { + doTestStrings(fetchSize, consistencyLevel, serialConsistencyLevel, executionProfile, cqlTemplateConsumer, null); + } + + private void doTestStrings(@Nullable Integer fetchSize, @Nullable ConsistencyLevel consistencyLevel, + @Nullable ConsistencyLevel serialConsistencyLevel, @Nullable String executionProfile, + Consumer cqlTemplateConsumer, @Nullable CqlIdentifier keyspace) { String[] results = { "Walter", "Hank", " Jesse" }; @@ -718,6 +737,10 @@ class CqlTemplateUnitTests { template.setExecutionProfile(executionProfile); } + if (keyspace != null) { + template.setKeyspace(keyspace); + } + cqlTemplateConsumer.accept(template); ArgumentCaptor statementArgumentCaptor = ArgumentCaptor.forClass(Statement.class); @@ -740,5 +763,8 @@ class CqlTemplateUnitTests { if (executionProfile != null) { assertThat(statement.getExecutionProfileName()).isEqualTo(executionProfile); } + if (keyspace != null) { + assertThat(statement.getKeyspace()).isEqualTo(keyspace); + } } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/QueryOptionsUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/QueryOptionsUnitTests.java index 2714904c6..e72c0e4e2 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/QueryOptionsUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/QueryOptionsUnitTests.java @@ -21,12 +21,14 @@ import java.time.Duration; import org.junit.jupiter.api.Test; +import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.DefaultConsistencyLevel; /** * Unit tests for {@link QueryOptions}. * * @author Mark Paluch + * @author Tomasz Lelek */ class QueryOptionsUnitTests { @@ -34,20 +36,21 @@ class QueryOptionsUnitTests { void buildQueryOptions() { QueryOptions queryOptions = QueryOptions.builder().consistencyLevel(DefaultConsistencyLevel.ANY) - .timeout(Duration.ofSeconds(1)).pageSize(10).tracing(true).build(); + .timeout(Duration.ofSeconds(1)).pageSize(10).tracing(true).keyspace(CqlIdentifier.fromCql("ks1")).build(); assertThat(queryOptions.getClass()).isEqualTo(QueryOptions.class); assertThat(queryOptions.getConsistencyLevel()).isEqualTo(DefaultConsistencyLevel.ANY); assertThat(queryOptions.getTimeout()).isEqualTo(Duration.ofSeconds(1)); assertThat(queryOptions.getPageSize()).isEqualTo(10); assertThat(queryOptions.getTracing()).isTrue(); + assertThat(queryOptions.getKeyspace()).isEqualTo(CqlIdentifier.fromCql("ks1")); } @Test // DATACASS-56 void buildQueryOptionsMutate() { QueryOptions queryOptions = QueryOptions.builder().consistencyLevel(DefaultConsistencyLevel.ANY) - .timeout(Duration.ofSeconds(1)).pageSize(10).tracing(true).build(); + .timeout(Duration.ofSeconds(1)).pageSize(10).tracing(true).keyspace(CqlIdentifier.fromCql("ks1")).build(); QueryOptions mutated = queryOptions.mutate().timeout(Duration.ofSeconds(5)).build(); @@ -58,5 +61,6 @@ class QueryOptionsUnitTests { assertThat(mutated.getTimeout()).isEqualTo(Duration.ofSeconds(5)); assertThat(mutated.getPageSize()).isEqualTo(10); assertThat(mutated.getTracing()).isTrue(); + assertThat(mutated.getKeyspace()).isEqualTo(CqlIdentifier.fromCql("ks1")); } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtilUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtilUnitTests.java index 687d87251..a822c186e 100755 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtilUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtilUnitTests.java @@ -15,6 +15,7 @@ */ package org.springframework.data.cassandra.core.cql; +import static org.assertj.core.api.Assertions.*; import static org.mockito.Mockito.*; import java.time.Duration; @@ -22,10 +23,14 @@ import java.util.concurrent.TimeUnit; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; + import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.DefaultConsistencyLevel; +import com.datastax.oss.driver.api.core.cql.BatchStatement; +import com.datastax.oss.driver.api.core.cql.BoundStatement; import com.datastax.oss.driver.api.core.cql.SimpleStatement; /** @@ -33,11 +38,14 @@ import com.datastax.oss.driver.api.core.cql.SimpleStatement; * * @author John Blum * @author Mark Paluch + * @author Tomasz Lelek */ @ExtendWith(MockitoExtension.class) class QueryOptionsUtilUnitTests { @Mock SimpleStatement simpleStatement; + @Mock BatchStatement batchStatement; + @Mock BoundStatement boundStatement; @Test // DATACASS-202, DATACASS-708 void addPreparedStatementOptionsShouldAddDriverQueryOptions() { @@ -88,4 +96,39 @@ class QueryOptionsUtilUnitTests { verify(simpleStatement).setPageSize(10); verify(simpleStatement).setTracing(true); } + + @Test // DATACASS-767 + void addKeyspaceOptionsOnSimpleStatementShouldAddDriverQueryOptions() { + when(simpleStatement.setKeyspace(any(CqlIdentifier.class))).thenReturn(simpleStatement); + + QueryOptions queryOptions = QueryOptions.builder() // + .keyspace(CqlIdentifier.fromCql("ks1")).build(); + + QueryOptionsUtil.addQueryOptions(simpleStatement, queryOptions); + + verify(simpleStatement).setKeyspace(CqlIdentifier.fromCql("ks1")); + } + + @Test // DATACASS-767 + void addKeyspaceOptionsOnBatchStatementShouldAddDriverQueryOptions() { + when(batchStatement.setKeyspace(any(CqlIdentifier.class))).thenReturn(batchStatement); + + QueryOptions queryOptions = QueryOptions.builder() // + .keyspace(CqlIdentifier.fromCql("ks1")).build(); + + QueryOptionsUtil.addQueryOptions(batchStatement, queryOptions); + + verify(batchStatement).setKeyspace(CqlIdentifier.fromCql("ks1")); + } + + @Test // DATACASS-767 + void addKeyspaceOptionsOnBoundStatementShouldThrowException() { + QueryOptions queryOptions = QueryOptions.builder() // + .keyspace(CqlIdentifier.fromCql("ks1")).build(); + + assertThatThrownBy(() -> QueryOptionsUtil.addQueryOptions(boundStatement, queryOptions)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Keyspace cannot be set for a BoundStatement"); + + } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplateIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplateIntegrationTests.java index 87899f039..1999beade 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplateIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplateIntegrationTests.java @@ -16,6 +16,7 @@ package org.springframework.data.cassandra.core.cql; import static org.assertj.core.api.Assertions.*; +import static org.junit.Assume.*; import reactor.test.StepVerifier; @@ -24,24 +25,30 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.springframework.data.cassandra.CassandraInvalidQueryException; import org.springframework.data.cassandra.ReactiveSession; import org.springframework.data.cassandra.core.cql.session.DefaultBridgedReactiveSession; import org.springframework.data.cassandra.core.cql.session.DefaultReactiveSessionFactory; +import org.springframework.data.cassandra.support.CassandraVersion; import org.springframework.data.cassandra.test.util.AbstractKeyspaceCreatingIntegrationTests; +import org.springframework.data.util.Version; +import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.cql.SimpleStatement; /** * Integration tests for {@link ReactiveCqlTemplate}. * * @author Mark Paluch + * @author Tomasz Lelek */ class ReactiveCqlTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrationTests { - + private static final Version CASSANDRA_4 = Version.parse("4.0"); private static final AtomicBoolean initialized = new AtomicBoolean(); private ReactiveSession reactiveSession; private ReactiveCqlTemplate template; + private Version cassandraVersion; @BeforeEach void before() { @@ -56,6 +63,7 @@ class ReactiveCqlTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegr getSession().execute("INSERT INTO user (id, username) VALUES ('WHITE', 'Walter');"); template = new ReactiveCqlTemplate(new DefaultReactiveSessionFactory(reactiveSession)); + cassandraVersion = CassandraVersion.get(getSession()); } @Test // DATACASS-335 @@ -138,4 +146,31 @@ class ReactiveCqlTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegr assertThat(actual).containsEntry("id", "WHITE").containsEntry("username", "Walter"); }).verifyComplete(); } + + @Test // DATACASS-767 + void selectByQueryWithKeyspaceShouldRetrieveData() { + assumeTrue(cassandraVersion.isGreaterThanOrEqualTo(CASSANDRA_4)); + + template.setKeyspace(CqlIdentifier.fromCql(keyspace)); + + template.queryForMap("SELECT * FROM user;").as(StepVerifier::create) // + .consumeNextWith(actual -> { + + assertThat(actual).containsEntry("id", "WHITE").containsEntry("username", "Walter"); + }).verifyComplete(); + } + + @Test // DATACASS-767 + void selectByQueryWithNonExistingKeyspaceShouldThrowThatKeyspaceDoesNotExists() { + assumeTrue(cassandraVersion.isGreaterThanOrEqualTo(CASSANDRA_4)); + + template.setKeyspace(CqlIdentifier.fromCql("non_existing")); + + template.queryForMap("SELECT * FROM user;").as(StepVerifier::create) // + .consumeErrorWith(e -> { + assertThat(e).isInstanceOf(CassandraInvalidQueryException.class) + .hasMessageContaining("Keyspace 'non_existing' does not exist"); + }).verify(); + + } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplateUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplateUnitTests.java index 8fe762fab..d8da69ea5 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplateUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplateUnitTests.java @@ -28,12 +28,12 @@ import java.util.function.Consumer; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; + import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.junit.jupiter.MockitoSettings; import org.mockito.quality.Strictness; - import org.springframework.dao.IncorrectResultSizeDataAccessException; import org.springframework.data.cassandra.CassandraConnectionFailureException; import org.springframework.data.cassandra.CassandraInvalidQueryException; @@ -44,6 +44,7 @@ import org.springframework.data.cassandra.core.cql.session.DefaultReactiveSessio import org.springframework.lang.Nullable; import com.datastax.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.DefaultConsistencyLevel; import com.datastax.oss.driver.api.core.NoNodeAvailableException; import com.datastax.oss.driver.api.core.cql.BoundStatement; @@ -58,6 +59,7 @@ import com.datastax.oss.driver.api.core.servererrors.InvalidQueryException; * Unit tests for {@link ReactiveCqlTemplate}. * * @author Mark Paluch + * @author Tomasz Lelek */ @ExtendWith(MockitoExtension.class) @MockitoSettings(strictness = Strictness.LENIENT) @@ -784,12 +786,18 @@ class ReactiveCqlTemplateUnitTests { } private void doTestStrings(Consumer cqlTemplateConsumer) { - doTestStrings(null, null, null, null, cqlTemplateConsumer); + doTestStrings(null, null, null, null, cqlTemplateConsumer, null); } private void doTestStrings(@Nullable Integer fetchSize, @Nullable ConsistencyLevel consistencyLevel, @Nullable ConsistencyLevel serialConsistencyLevel, @Nullable String executionProfile, Consumer cqlTemplateConsumer) { + doTestStrings(fetchSize, consistencyLevel, serialConsistencyLevel, executionProfile, cqlTemplateConsumer, null); + } + + private void doTestStrings(@Nullable Integer fetchSize, @Nullable ConsistencyLevel consistencyLevel, + @Nullable ConsistencyLevel serialConsistencyLevel, @Nullable String executionProfile, + Consumer cqlTemplateConsumer, @Nullable CqlIdentifier keyspace) { String[] results = { "Walter", "Hank", " Jesse" }; @@ -818,6 +826,10 @@ class ReactiveCqlTemplateUnitTests { template.setExecutionProfile(executionProfile); } + if (keyspace != null) { + template.setKeyspace(keyspace); + } + cqlTemplateConsumer.accept(template); ArgumentCaptor statementArgumentCaptor = ArgumentCaptor.forClass(Statement.class); @@ -840,5 +852,8 @@ class ReactiveCqlTemplateUnitTests { if (executionProfile != null) { assertThat(statement.getExecutionProfileName()).isEqualTo(executionProfile); } + if (keyspace != null) { + assertThat(statement.getKeyspace()).isEqualTo(keyspace); + } } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/test/util/CassandraDelegate.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/test/util/CassandraDelegate.java index e1903c107..69e28120f 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/test/util/CassandraDelegate.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/test/util/CassandraDelegate.java @@ -28,12 +28,14 @@ import org.springframework.data.cassandra.support.CqlDataSet; import org.springframework.util.Assert; import org.springframework.util.SocketUtils; import org.springframework.util.StringUtils; - import org.testcontainers.containers.CassandraContainer; import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.CqlSessionBuilder; +import com.datastax.oss.driver.api.core.config.DefaultDriverOption; +import com.datastax.oss.driver.api.core.config.DriverConfigLoader; +import com.datastax.oss.driver.api.core.metadata.Node; /** * Delegate used to provide a Cassandra context for integration tests. This rule can use/spin up either an embedded @@ -43,6 +45,7 @@ import com.datastax.oss.driver.api.core.CqlSessionBuilder; * * @author Mark Paluch * @author John Blum + * @author Tomasz Lelek * @since 1.5 * @see CassandraConnectionProperties */ @@ -296,8 +299,21 @@ class CassandraDelegate { String host = resolveHost(); - return CqlSession.builder().addContactPoint(InetSocketAddress.createUnresolved(host, port)) + CqlSessionBuilder builder = CqlSession.builder().addContactPoint(InetSocketAddress.createUnresolved(host, port)) .withLocalDatacenter("datacenter1"); + + CqlSession cqlSession = builder.build(); + if (cassandraVersionGreaterThanOrEqualTo4(cqlSession)) { + return builder.withConfigLoader( + DriverConfigLoader.programmaticBuilder().withString(DefaultDriverOption.PROTOCOL_VERSION, "V5").build()); + } else { + return builder; + } + } + + private boolean cassandraVersionGreaterThanOrEqualTo4(CqlSession cqlSession) { + return cqlSession.getMetadata().getNodes().values().stream().map(Node::getCassandraVersion) + .allMatch(v -> v != null && v.getMajor() >= 4); } private String resolveHost() {