@@ -17,6 +17,7 @@ package org.springframework.data.cassandra.core;
|
||||
|
||||
import java.util.Collections;
|
||||
|
||||
import org.springframework.data.cassandra.core.cql.QueryOptions;
|
||||
import org.springframework.data.cassandra.core.cql.WriteOptions;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
@@ -59,6 +60,16 @@ public interface CassandraBatchOperations {
|
||||
*/
|
||||
CassandraBatchOperations withTimestamp(long timestamp);
|
||||
|
||||
/**
|
||||
* Apply given {@link QueryOptions} to the whole batch statement.
|
||||
*
|
||||
* @param options the options to apply.
|
||||
* @return {@code this} {@link CassandraBatchOperations}.
|
||||
* @throws IllegalStateException if the batch was already executed.
|
||||
* @since 4.4
|
||||
*/
|
||||
CassandraBatchOperations withQueryOptions(QueryOptions options);
|
||||
|
||||
/**
|
||||
* Add a {@link BatchableStatement statement} to the batch.
|
||||
*
|
||||
|
||||
@@ -20,6 +20,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import org.springframework.data.cassandra.core.convert.CassandraConverter;
|
||||
import org.springframework.data.cassandra.core.cql.QueryOptions;
|
||||
import org.springframework.data.cassandra.core.cql.QueryOptionsUtil;
|
||||
import org.springframework.data.cassandra.core.cql.WriteOptions;
|
||||
import org.springframework.data.cassandra.core.mapping.BasicCassandraPersistentEntity;
|
||||
import org.springframework.data.cassandra.core.mapping.CassandraMappingContext;
|
||||
@@ -58,6 +59,8 @@ class CassandraBatchTemplate implements CassandraBatchOperations {
|
||||
|
||||
private final StatementFactory statementFactory;
|
||||
|
||||
private QueryOptions options = QueryOptions.empty();
|
||||
|
||||
/**
|
||||
* Create a new {@link CassandraBatchTemplate} given {@link CassandraOperations} and {@link BatchType}.
|
||||
*
|
||||
@@ -114,7 +117,8 @@ class CassandraBatchTemplate implements CassandraBatchOperations {
|
||||
public WriteResult execute() {
|
||||
|
||||
if (this.executed.compareAndSet(false, true)) {
|
||||
return WriteResult.of(this.operations.getCqlOperations().queryForResultSet(batch.build()));
|
||||
BatchStatement statement = QueryOptionsUtil.addQueryOptions(batch.build(), this.options);
|
||||
return WriteResult.of(this.operations.getCqlOperations().queryForResultSet(statement));
|
||||
}
|
||||
|
||||
throw new IllegalStateException("This Cassandra Batch was already executed");
|
||||
@@ -130,9 +134,21 @@ class CassandraBatchTemplate implements CassandraBatchOperations {
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CassandraBatchOperations withQueryOptions(QueryOptions options) {
|
||||
|
||||
assertNotExecuted();
|
||||
Assert.notNull(options, "QueryOptions must not be null");
|
||||
|
||||
this.options = options;
|
||||
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CassandraBatchOperations addStatement(BatchableStatement<?> statement) {
|
||||
|
||||
assertNotExecuted();
|
||||
Assert.notNull(statement, "Statement must not be null");
|
||||
|
||||
this.batch.addStatement(statement);
|
||||
@@ -143,6 +159,7 @@ class CassandraBatchTemplate implements CassandraBatchOperations {
|
||||
@Override
|
||||
public CassandraBatchOperations addStatements(BatchableStatement<?>... statements) {
|
||||
|
||||
assertNotExecuted();
|
||||
Assert.notNull(statements, "Statements must not be null");
|
||||
|
||||
this.batch.addStatements(statements);
|
||||
@@ -154,6 +171,7 @@ class CassandraBatchTemplate implements CassandraBatchOperations {
|
||||
@SuppressWarnings("unchecked")
|
||||
public CassandraBatchOperations addStatements(Iterable<? extends BatchableStatement<?>> statements) {
|
||||
|
||||
assertNotExecuted();
|
||||
Assert.notNull(statements, "Statements must not be null");
|
||||
|
||||
this.batch.addStatements((Iterable<BatchableStatement<?>>) statements);
|
||||
|
||||
@@ -22,6 +22,7 @@ import java.util.Collections;
|
||||
|
||||
import org.reactivestreams.Subscriber;
|
||||
|
||||
import org.springframework.data.cassandra.core.cql.QueryOptions;
|
||||
import org.springframework.data.cassandra.core.cql.WriteOptions;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
@@ -65,6 +66,16 @@ public interface ReactiveCassandraBatchOperations {
|
||||
*/
|
||||
ReactiveCassandraBatchOperations withTimestamp(long timestamp);
|
||||
|
||||
/**
|
||||
* Apply given {@link QueryOptions} to the whole batch statement.
|
||||
*
|
||||
* @param options the options to apply.
|
||||
* @return {@code this} {@link CassandraBatchOperations}.
|
||||
* @throws IllegalStateException if the batch was already executed.
|
||||
* @since 4.4
|
||||
*/
|
||||
ReactiveCassandraBatchOperations withQueryOptions(QueryOptions options);
|
||||
|
||||
/**
|
||||
* Add a {@link BatchableStatement statement} to the batch.
|
||||
*
|
||||
|
||||
@@ -27,6 +27,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import org.springframework.data.cassandra.core.convert.CassandraConverter;
|
||||
import org.springframework.data.cassandra.core.cql.QueryOptions;
|
||||
import org.springframework.data.cassandra.core.cql.QueryOptionsUtil;
|
||||
import org.springframework.data.cassandra.core.cql.WriteOptions;
|
||||
import org.springframework.data.cassandra.core.mapping.BasicCassandraPersistentEntity;
|
||||
import org.springframework.data.cassandra.core.mapping.CassandraMappingContext;
|
||||
@@ -66,6 +67,8 @@ class ReactiveCassandraBatchTemplate implements ReactiveCassandraBatchOperations
|
||||
|
||||
private final StatementFactory statementFactory;
|
||||
|
||||
private QueryOptions options = QueryOptions.empty();
|
||||
|
||||
/**
|
||||
* Create a new {@link CassandraBatchTemplate} given {@link CassandraOperations} and {@link BatchType}.
|
||||
*
|
||||
@@ -141,7 +144,8 @@ class ReactiveCassandraBatchTemplate implements ReactiveCassandraBatchOperations
|
||||
|
||||
this.batch.addStatements((List<BatchableStatement<?>>) statements);
|
||||
|
||||
return this.operations.getReactiveCqlOperations().queryForResultSet(this.batch.build());
|
||||
return this.operations.getReactiveCqlOperations()
|
||||
.queryForResultSet(QueryOptionsUtil.addQueryOptions(this.batch.build(), this.options));
|
||||
}) //
|
||||
.flatMap(resultSet -> resultSet.rows().collectList()
|
||||
.map(rows -> new WriteResult(resultSet.getAllExecutionInfo(), resultSet.wasApplied(), rows)));
|
||||
@@ -160,9 +164,21 @@ class ReactiveCassandraBatchTemplate implements ReactiveCassandraBatchOperations
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ReactiveCassandraBatchOperations withQueryOptions(QueryOptions options) {
|
||||
|
||||
assertNotExecuted();
|
||||
Assert.notNull(options, "QueryOptions must not be null");
|
||||
|
||||
this.options = options;
|
||||
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ReactiveCassandraBatchOperations addStatement(Mono<? extends BatchableStatement<?>> statement) {
|
||||
|
||||
assertNotExecuted();
|
||||
Assert.notNull(statement, "Statement mono must not be null");
|
||||
|
||||
this.batchMonos.add(statement.map(List::of));
|
||||
@@ -174,6 +190,7 @@ class ReactiveCassandraBatchTemplate implements ReactiveCassandraBatchOperations
|
||||
public ReactiveCassandraBatchOperations addStatements(
|
||||
Mono<? extends Iterable<? extends BatchableStatement<?>>> statements) {
|
||||
|
||||
assertNotExecuted();
|
||||
Assert.notNull(statements, "Statements mono must not be null");
|
||||
|
||||
this.batchMonos.add(statements);
|
||||
|
||||
@@ -24,6 +24,8 @@ import java.util.concurrent.TimeUnit;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import org.springframework.data.cassandra.CassandraConnectionFailureException;
|
||||
import org.springframework.data.cassandra.core.cql.QueryOptions;
|
||||
import org.springframework.data.cassandra.core.cql.WriteOptions;
|
||||
import org.springframework.data.cassandra.domain.FlatGroup;
|
||||
import org.springframework.data.cassandra.domain.Group;
|
||||
@@ -31,6 +33,8 @@ import org.springframework.data.cassandra.domain.GroupKey;
|
||||
import org.springframework.data.cassandra.repository.support.SchemaTestUtils;
|
||||
import org.springframework.data.cassandra.test.util.AbstractKeyspaceCreatingIntegrationTests;
|
||||
|
||||
import com.datastax.oss.driver.api.core.AllNodesFailedException;
|
||||
import com.datastax.oss.driver.api.core.ConsistencyLevel;
|
||||
import com.datastax.oss.driver.api.core.cql.BatchType;
|
||||
import com.datastax.oss.driver.api.core.cql.ResultSet;
|
||||
import com.datastax.oss.driver.api.core.cql.Row;
|
||||
@@ -297,7 +301,19 @@ class CassandraBatchTemplateIntegrationTests extends AbstractKeyspaceCreatingInt
|
||||
}
|
||||
}
|
||||
|
||||
@Test // DATACASS-288
|
||||
@Test // GH-1192
|
||||
void shouldApplyQueryOptions() {
|
||||
|
||||
QueryOptions options = QueryOptions.builder().consistencyLevel(ConsistencyLevel.THREE).build();
|
||||
|
||||
CassandraBatchOperations batchOperations = new CassandraBatchTemplate(template, BatchType.LOGGED);
|
||||
CassandraBatchOperations ops = batchOperations.insert(walter).withQueryOptions(options);
|
||||
|
||||
assertThatExceptionOfType(CassandraConnectionFailureException.class).isThrownBy(ops::execute)
|
||||
.withRootCauseInstanceOf(AllNodesFailedException.class);
|
||||
}
|
||||
|
||||
@Test // GH-1192
|
||||
void shouldNotExecuteTwice() {
|
||||
|
||||
CassandraBatchOperations batchOperations = new CassandraBatchTemplate(template, BatchType.LOGGED);
|
||||
|
||||
@@ -30,8 +30,10 @@ import java.util.concurrent.TimeUnit;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import org.springframework.data.cassandra.CassandraConnectionFailureException;
|
||||
import org.springframework.data.cassandra.ReactiveResultSet;
|
||||
import org.springframework.data.cassandra.core.convert.MappingCassandraConverter;
|
||||
import org.springframework.data.cassandra.core.cql.QueryOptions;
|
||||
import org.springframework.data.cassandra.core.cql.ReactiveCqlTemplate;
|
||||
import org.springframework.data.cassandra.core.cql.WriteOptions;
|
||||
import org.springframework.data.cassandra.core.cql.session.DefaultBridgedReactiveSession;
|
||||
@@ -41,6 +43,8 @@ import org.springframework.data.cassandra.domain.GroupKey;
|
||||
import org.springframework.data.cassandra.repository.support.SchemaTestUtils;
|
||||
import org.springframework.data.cassandra.test.util.AbstractKeyspaceCreatingIntegrationTests;
|
||||
|
||||
import com.datastax.oss.driver.api.core.AllNodesFailedException;
|
||||
import com.datastax.oss.driver.api.core.ConsistencyLevel;
|
||||
import com.datastax.oss.driver.api.core.cql.BatchType;
|
||||
import com.datastax.oss.driver.api.core.cql.Row;
|
||||
import com.datastax.oss.driver.api.core.cql.SimpleStatement;
|
||||
@@ -153,8 +157,7 @@ class ReactiveCassandraBatchTemplateIntegrationTests extends AbstractKeyspaceCre
|
||||
.then(template.getReactiveCqlOperations().queryForResultSet("SELECT TTL(email), email FROM group;"));
|
||||
|
||||
resultSet.flatMapMany(ReactiveResultSet::availableRows) //
|
||||
.collectList()
|
||||
.as(StepVerifier::create) //
|
||||
.collectList().as(StepVerifier::create) //
|
||||
.assertNext(rows -> {
|
||||
|
||||
for (Row row : rows) {
|
||||
@@ -252,8 +255,7 @@ class ReactiveCassandraBatchTemplateIntegrationTests extends AbstractKeyspaceCre
|
||||
.then(template.getReactiveCqlOperations().queryForResultSet("SELECT TTL(email), email FROM group;"));
|
||||
|
||||
resultSet.flatMapMany(ReactiveResultSet::availableRows) //
|
||||
.collectList()
|
||||
.as(StepVerifier::create) //
|
||||
.collectList().as(StepVerifier::create) //
|
||||
.assertNext(rows -> {
|
||||
|
||||
for (Row row : rows) {
|
||||
@@ -388,6 +390,20 @@ class ReactiveCassandraBatchTemplateIntegrationTests extends AbstractKeyspaceCre
|
||||
.assertNext(row -> assertThat(row.getLong(0)).isEqualTo(timestamp)).verifyComplete();
|
||||
}
|
||||
|
||||
@Test // GH-1192
|
||||
void shouldApplyQueryOptions() {
|
||||
|
||||
QueryOptions options = QueryOptions.builder().consistencyLevel(ConsistencyLevel.THREE).build();
|
||||
|
||||
ReactiveCassandraBatchOperations batchOperations = new ReactiveCassandraBatchTemplate(template, BatchType.LOGGED);
|
||||
Mono<WriteResult> execute = batchOperations.insert(walter).insert(mike).withQueryOptions(options).execute();
|
||||
|
||||
execute.as(StepVerifier::create).verifyErrorSatisfies(e -> {
|
||||
assertThat(e).isInstanceOf(CassandraConnectionFailureException.class)
|
||||
.hasRootCauseInstanceOf(AllNodesFailedException.class);
|
||||
});
|
||||
}
|
||||
|
||||
@Test // DATACASS-574
|
||||
void shouldNotExecuteTwice() {
|
||||
|
||||
@@ -428,10 +444,12 @@ class ReactiveCassandraBatchTemplateIntegrationTests extends AbstractKeyspaceCre
|
||||
|
||||
for (int i = 0; i < 100; i++) {
|
||||
|
||||
batchOperations.insert(Mono.just(Arrays.asList(new Group(new GroupKey("users", "0x1", "walter" + random.longs())),
|
||||
new Group(new GroupKey("users", "0x1", "walter" + random.longs())),
|
||||
new Group(new GroupKey("users", "0x1", "walter" + random.longs())),
|
||||
new Group(new GroupKey("users", "0x1", "walter" + random.longs())))).publishOn(Schedulers.boundedElastic()));
|
||||
batchOperations.insert(Mono
|
||||
.just(Arrays.asList(new Group(new GroupKey("users", "0x1", "walter" + random.longs())),
|
||||
new Group(new GroupKey("users", "0x1", "walter" + random.longs())),
|
||||
new Group(new GroupKey("users", "0x1", "walter" + random.longs())),
|
||||
new Group(new GroupKey("users", "0x1", "walter" + random.longs()))))
|
||||
.publishOn(Schedulers.boundedElastic()));
|
||||
}
|
||||
|
||||
batchOperations.execute()
|
||||
|
||||
Reference in New Issue
Block a user