DATACASS-611 - Add query derivation for delete queries.
We now support query derivation for delete queries using delete…By in method declarations.
interface PersonRepository extends Repository<Person, String> {
void deleteWithoutResultByLastname(String lastname);
boolean deleteByLastname(String lastname);
}
This commit is contained in:
@@ -146,6 +146,8 @@ public abstract class AbstractCassandraQuery extends CassandraRepositoryQuerySup
|
||||
return ((statement, type) -> new SingleEntityExecution(getOperations(), false).execute(statement, Long.class));
|
||||
} else if (isExistsQuery()) {
|
||||
return new ExistsExecution(getOperations());
|
||||
} else if (isModifyingQuery()) {
|
||||
return ((statement, type) -> getOperations().getCqlOperations().queryForResultSet(statement).wasApplied());
|
||||
} else {
|
||||
return new SingleEntityExecution(getOperations(), isLimiting());
|
||||
}
|
||||
@@ -174,4 +176,12 @@ public abstract class AbstractCassandraQuery extends CassandraRepositoryQuerySup
|
||||
* @since 2.0.4
|
||||
*/
|
||||
protected abstract boolean isLimiting();
|
||||
|
||||
/**
|
||||
* Returns whether the query is a modifying query.
|
||||
*
|
||||
* @return a boolean value indicating whether the query is a modifying query.
|
||||
* @since 2.2
|
||||
*/
|
||||
protected abstract boolean isModifyingQuery();
|
||||
}
|
||||
|
||||
@@ -20,6 +20,7 @@ import reactor.core.publisher.Mono;
|
||||
|
||||
import org.reactivestreams.Publisher;
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.data.cassandra.ReactiveResultSet;
|
||||
import org.springframework.data.cassandra.core.CassandraOperations;
|
||||
import org.springframework.data.cassandra.core.ReactiveCassandraOperations;
|
||||
import org.springframework.data.cassandra.core.convert.CassandraConverter;
|
||||
@@ -151,6 +152,9 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository
|
||||
Long.class));
|
||||
} else if (isExistsQuery()) {
|
||||
return new ExistsExecution(getReactiveCassandraOperations());
|
||||
} else if (isModifyingQuery()) {
|
||||
return (statement, type) -> getReactiveCassandraOperations().getReactiveCqlOperations()
|
||||
.queryForResultSet(statement).map(ReactiveResultSet::wasApplied);
|
||||
} else {
|
||||
return new SingleEntityExecution(getReactiveCassandraOperations(), isLimiting());
|
||||
}
|
||||
@@ -180,6 +184,14 @@ public abstract class AbstractReactiveCassandraQuery extends CassandraRepository
|
||||
*/
|
||||
protected abstract boolean isLimiting();
|
||||
|
||||
/**
|
||||
* Returns whether the query is a modifying query.
|
||||
*
|
||||
* @return a boolean value indicating whether the query is a modifying query.
|
||||
* @since 2.2
|
||||
*/
|
||||
protected abstract boolean isModifyingQuery();
|
||||
|
||||
private static CassandraConverter getRequiredConverter(ReactiveCassandraOperations operations) {
|
||||
|
||||
Assert.notNull(operations, "ReactiveCassandraOperations must not be null");
|
||||
|
||||
@@ -104,6 +104,10 @@ public class PartTreeCassandraQuery extends AbstractCassandraQuery {
|
||||
return getQueryStatementCreator().exists(getStatementFactory(), getTree(), parameterAccessor);
|
||||
}
|
||||
|
||||
if (getTree().isDelete()) {
|
||||
return getQueryStatementCreator().delete(getStatementFactory(), getTree(), parameterAccessor);
|
||||
}
|
||||
|
||||
return getQueryStatementCreator().select(getStatementFactory(), getTree(), parameterAccessor,
|
||||
getQueryMethod().getResultProcessor());
|
||||
}
|
||||
@@ -131,4 +135,12 @@ public class PartTreeCassandraQuery extends AbstractCassandraQuery {
|
||||
protected boolean isLimiting() {
|
||||
return getTree().isLimiting();
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.data.cassandra.repository.query.AbstractCassandraQuery#isModifyingQuery()
|
||||
*/
|
||||
@Override
|
||||
protected boolean isModifyingQuery() {
|
||||
return getTree().isDelete();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -118,6 +118,32 @@ class QueryStatementCreator {
|
||||
return doWithQuery(parameterAccessor, tree, function);
|
||||
}
|
||||
|
||||
/**
|
||||
* Create a {@literal DELETE} {@link Statement} from a {@link PartTree} and apply query options for delete query
|
||||
* execution.
|
||||
*
|
||||
* @param statementFactory must not be {@literal null}.
|
||||
* @param tree must not be {@literal null}.
|
||||
* @param parameterAccessor must not be {@literal null}.
|
||||
* @return the {@literal DELETE} {@link Statement}.
|
||||
* @since 2.2
|
||||
*/
|
||||
Statement delete(StatementFactory statementFactory, PartTree tree, CassandraParameterAccessor parameterAccessor) {
|
||||
|
||||
Function<Query, Statement> function = query -> {
|
||||
|
||||
RegularStatement statement = statementFactory.delete(query, requirePersistentEntity());
|
||||
|
||||
if (LOG.isDebugEnabled()) {
|
||||
LOG.debug(String.format("Created query [%s].", statement));
|
||||
}
|
||||
|
||||
return statement;
|
||||
};
|
||||
|
||||
return doWithQuery(parameterAccessor, tree, function);
|
||||
}
|
||||
|
||||
/**
|
||||
* Create a {@literal SELECT} {@link Statement} from a {@link PartTree} and apply query options for exists query
|
||||
* execution. Limit results to a single row.
|
||||
|
||||
@@ -17,10 +17,13 @@ package org.springframework.data.cassandra.repository.query;
|
||||
|
||||
import lombok.NonNull;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
import org.reactivestreams.Publisher;
|
||||
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.dao.IncorrectResultSizeDataAccessException;
|
||||
import org.springframework.data.cassandra.core.ReactiveCassandraOperations;
|
||||
@@ -190,7 +193,6 @@ interface ReactiveCassandraQueryExecution {
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.data.cassandra.repository.query.ReactiveCassandraQueryExecution#execute(java.lang.String, java.lang.Class)
|
||||
*/
|
||||
@SuppressWarnings("ConstantConditions")
|
||||
@Override
|
||||
public Object execute(Statement statement, Class<?> type) {
|
||||
return converter.convert(delegate.execute(statement, type));
|
||||
@@ -222,6 +224,16 @@ interface ReactiveCassandraQueryExecution {
|
||||
return source;
|
||||
}
|
||||
|
||||
if (returnedType.getReturnedType().equals(Void.class)) {
|
||||
if (source instanceof Mono) {
|
||||
return ((Mono<?>) source).then();
|
||||
}
|
||||
|
||||
if (source instanceof Publisher) {
|
||||
return Flux.from((Publisher<?>) source).then();
|
||||
}
|
||||
}
|
||||
|
||||
if (returnedType.isInstance(source)) {
|
||||
return source;
|
||||
}
|
||||
|
||||
@@ -104,6 +104,10 @@ public class ReactivePartTreeCassandraQuery extends AbstractReactiveCassandraQue
|
||||
return getQueryStatementCreator().exists(getStatementFactory(), getTree(), parameterAccessor);
|
||||
}
|
||||
|
||||
if (getTree().isDelete()) {
|
||||
return getQueryStatementCreator().delete(getStatementFactory(), getTree(), parameterAccessor);
|
||||
}
|
||||
|
||||
return getQueryStatementCreator().select(getStatementFactory(), getTree(), parameterAccessor,
|
||||
getQueryMethod().getResultProcessor());
|
||||
}
|
||||
@@ -131,4 +135,12 @@ public class ReactivePartTreeCassandraQuery extends AbstractReactiveCassandraQue
|
||||
protected boolean isLimiting() {
|
||||
return getTree().isLimiting();
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.data.cassandra.repository.query.AbstractCassandraQuery#isModifyingQuery()
|
||||
*/
|
||||
@Override
|
||||
protected boolean isModifyingQuery() {
|
||||
return getTree().isDelete();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -140,4 +140,12 @@ public class ReactiveStringBasedCassandraQuery extends AbstractReactiveCassandra
|
||||
protected boolean isLimiting() {
|
||||
return false;
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.data.cassandra.repository.query.AbstractCassandraQuery#isModifyingQuery()
|
||||
*/
|
||||
@Override
|
||||
protected boolean isModifyingQuery() {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -134,4 +134,12 @@ public class StringBasedCassandraQuery extends AbstractCassandraQuery {
|
||||
protected boolean isLimiting() {
|
||||
return false;
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.data.cassandra.repository.query.AbstractCassandraQuery#isModifyingQuery()
|
||||
*/
|
||||
@Override
|
||||
protected boolean isModifyingQuery() {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -343,6 +343,23 @@ public class QueryDerivationIntegrationTests extends AbstractSpringDataEmbeddedC
|
||||
assertThat(count).isEqualTo(3);
|
||||
}
|
||||
|
||||
@Test // DATACASS-611
|
||||
public void shouldDeleteRecords() {
|
||||
|
||||
personRepository.deleteByLastname("White");
|
||||
|
||||
assertThat(personRepository.countByLastname("White")).isZero();
|
||||
}
|
||||
|
||||
@Test // DATACASS-611
|
||||
public void shouldDeleteRecordsWithWasApplied() {
|
||||
|
||||
boolean deleted = personRepository.deleteByLastname("White");
|
||||
|
||||
assertThat(deleted).isTrue();
|
||||
assertThat(personRepository.countByLastname("White")).isZero();
|
||||
}
|
||||
|
||||
@Test // DATACASS-512
|
||||
public void shouldApplyExistsProjection() {
|
||||
|
||||
@@ -383,6 +400,10 @@ public class QueryDerivationIntegrationTests extends AbstractSpringDataEmbeddedC
|
||||
|
||||
long countByLastname(String lastname);
|
||||
|
||||
boolean deleteByLastname(String lastname);
|
||||
|
||||
void deleteVoidByLastname(String lastname);
|
||||
|
||||
boolean existsByLastname(String lastname);
|
||||
|
||||
Slice<Person> findAllSlicedByLastname(String lastname, Pageable pageable);
|
||||
|
||||
@@ -15,19 +15,19 @@
|
||||
*/
|
||||
package org.springframework.data.cassandra.repository;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.HashSet;
|
||||
import java.util.Set;
|
||||
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
import reactor.test.StepVerifier;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.HashSet;
|
||||
import java.util.Set;
|
||||
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
|
||||
import org.reactivestreams.Publisher;
|
||||
|
||||
import org.springframework.beans.BeansException;
|
||||
import org.springframework.beans.factory.BeanClassLoaderAware;
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
@@ -132,11 +132,10 @@ public class ReactiveCassandraRepositoryIntegrationTests extends AbstractKeyspac
|
||||
repository.findByLastname(dave.getLastname()).as(StepVerifier::create).expectNextCount(2).verifyComplete();
|
||||
}
|
||||
|
||||
@Test //DATACASS-529
|
||||
@Test // DATACASS-529
|
||||
public void shouldFindSliceByLastName() {
|
||||
repository.findByLastname(carter.getLastname(), CassandraPageRequest.first(1)).as(StepVerifier::create)
|
||||
.expectNextMatches(users -> users.getSize() == 1 && users.hasNext())
|
||||
.verifyComplete();
|
||||
.expectNextMatches(users -> users.getSize() == 1 && users.hasNext()).verifyComplete();
|
||||
}
|
||||
|
||||
@Test // DATACASS-529
|
||||
@@ -193,8 +192,7 @@ public class ReactiveCassandraRepositoryIntegrationTests extends AbstractKeyspac
|
||||
.verifyComplete();
|
||||
|
||||
groupRepostitory.findByIdGroupnameAndIdHashPrefix("Simpsons", "hash", Sort.by("id.username").descending())
|
||||
.as(StepVerifier::create)
|
||||
.expectNext(new Group(key2), new Group(key1)) //
|
||||
.as(StepVerifier::create).expectNext(new Group(key2), new Group(key1)) //
|
||||
.verifyComplete();
|
||||
}
|
||||
|
||||
@@ -208,6 +206,22 @@ public class ReactiveCassandraRepositoryIntegrationTests extends AbstractKeyspac
|
||||
repository.countQueryByLastname("None").as(StepVerifier::create).expectNext(0L).verifyComplete();
|
||||
}
|
||||
|
||||
@Test // DATACASS-611
|
||||
public void shouldDeleteRecords() {
|
||||
|
||||
repository.deleteVoidById(dave.getId()).as(StepVerifier::create).verifyComplete();
|
||||
|
||||
repository.countByLastname("Matthews").as(StepVerifier::create).expectNext(1L).verifyComplete();
|
||||
}
|
||||
|
||||
@Test // DATACASS-611
|
||||
public void shouldDeleteRecordsReturingWasApplied() {
|
||||
|
||||
repository.deleteAllById(dave.getId()).as(StepVerifier::create).expectNext(true).verifyComplete();
|
||||
|
||||
repository.countByLastname("Matthews").as(StepVerifier::create).expectNext(1L).verifyComplete();
|
||||
}
|
||||
|
||||
@Test // DATACASS-512
|
||||
public void shouldApplyExistsProjection() {
|
||||
|
||||
@@ -232,6 +246,10 @@ public class ReactiveCassandraRepositoryIntegrationTests extends AbstractKeyspac
|
||||
|
||||
Mono<Long> countByLastname(String lastname);
|
||||
|
||||
Mono<Boolean> deleteAllById(String lastname);
|
||||
|
||||
Mono<Void> deleteVoidById(String lastname);
|
||||
|
||||
Mono<Boolean> existsByLastname(String lastname);
|
||||
|
||||
@Query("SELECT * FROM users WHERE lastname = ?0")
|
||||
|
||||
@@ -30,6 +30,7 @@ import org.junit.rules.ExpectedException;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.mockito.Mock;
|
||||
import org.mockito.junit.MockitoJUnitRunner;
|
||||
|
||||
import org.springframework.data.cassandra.core.CassandraOperations;
|
||||
import org.springframework.data.cassandra.core.convert.CassandraConverter;
|
||||
import org.springframework.data.cassandra.core.convert.MappingCassandraConverter;
|
||||
@@ -215,6 +216,15 @@ public class PartTreeCassandraQueryUnitTests {
|
||||
assertThat(statement.toString()).isEqualTo("SELECT COUNT(1) FROM person;");
|
||||
}
|
||||
|
||||
@Test // DATACASS-611
|
||||
public void shouldCreateDeleteQuery() {
|
||||
|
||||
Statement statement = deriveQueryFromMethod(Repo.class, "deleteAllByLastname", new Class[] { String.class },
|
||||
"Walter");
|
||||
|
||||
assertThat(statement.toString()).isEqualTo("DELETE FROM person WHERE lastname='Walter';");
|
||||
}
|
||||
|
||||
@Test // DATACASS-512
|
||||
public void shouldCreateExistsQuery() {
|
||||
|
||||
@@ -295,6 +305,8 @@ public class PartTreeCassandraQueryUnitTests {
|
||||
|
||||
long countBy();
|
||||
|
||||
boolean deleteAllByLastname(String lastname);
|
||||
|
||||
boolean existsBy();
|
||||
|
||||
@AllowFiltering
|
||||
|
||||
@@ -141,6 +141,15 @@ public class ReactivePartTreeCassandraQueryUnitTests {
|
||||
assertThat(statement.toString()).isEqualTo("SELECT COUNT(1) FROM person;");
|
||||
}
|
||||
|
||||
@Test // DATACASS-611
|
||||
public void shouldCreateDeleteQuery() {
|
||||
|
||||
Statement statement = deriveQueryFromMethod(PartTreeCassandraQueryUnitTests.Repo.class, "deleteAllByLastname",
|
||||
new Class[] { String.class }, "Walter");
|
||||
|
||||
assertThat(statement.toString()).isEqualTo("DELETE FROM person WHERE lastname='Walter';");
|
||||
}
|
||||
|
||||
@Test // DATACASS-512
|
||||
public void shouldCreateExistsQuery() {
|
||||
|
||||
@@ -202,6 +211,8 @@ public class ReactivePartTreeCassandraQueryUnitTests {
|
||||
|
||||
Mono<Long> countBy();
|
||||
|
||||
Mono<Boolean> deleteAllByLastname(String lastname);
|
||||
|
||||
Mono<Boolean> existsBy();
|
||||
|
||||
@Consistency(ConsistencyLevel.LOCAL_ONE)
|
||||
|
||||
@@ -13,6 +13,7 @@ This chapter summarizes changes and new features for each release.
|
||||
* Optimistic Locking support.
|
||||
* Auditing via `@EnableCassandraAuditing`.
|
||||
* Idempotency support in `@Query` annotation.
|
||||
* Query derivation for <<cassandra.repositories.queries.delete,`DELETE` queries>>.
|
||||
|
||||
[[new-features.2-1-0]]
|
||||
== What's new in Spring Data for Apache Cassandra 2.1
|
||||
|
||||
@@ -282,6 +282,25 @@ lower / upper bounds (`>` / `>=` & `<` / `<=`) according to `Range`
|
||||
|
||||
|===
|
||||
|
||||
[[cassandra.repositories.queries.delete]]
|
||||
== Repository Delete Queries
|
||||
|
||||
The keywords in the preceding table can be used in conjunction with `delete…By` to create queries that delete matching documents.
|
||||
|
||||
====
|
||||
[source,java]
|
||||
----
|
||||
interface PersonRepository extends Repository<Person, String> {
|
||||
|
||||
void deleteWithoutResultByLastname(String lastname);
|
||||
|
||||
boolean deleteByLastname(String lastname);
|
||||
}
|
||||
----
|
||||
====
|
||||
|
||||
Delete queries return whether the query was applied or terminate without returning a value using `void`.
|
||||
|
||||
include::../{spring-data-commons-docs}/repository-projections.adoc[leveloffset=+2]
|
||||
|
||||
[[cassandra.repositories.queries.options]]
|
||||
@@ -297,6 +316,7 @@ The declared consistency level is applied to the query each time it is executed.
|
||||
The following example sets the consistency level to `ConsistencyLevel.LOCAL_ONE`:
|
||||
|
||||
====
|
||||
[source,java]
|
||||
----
|
||||
public interface PersonRepository extends CrudRepository<Person, String> {
|
||||
|
||||
|
||||
Reference in New Issue
Block a user