Use queryScanConsistency on reactive deleteAll().

It was present on non-Reactive, but missing from reactive.

Closes #1096.
Original pull request: #1108.

Co-authored-by: mikereiche <michael.reiche@couchbase.com>
This commit is contained in:
Michael Reiche
2021-03-24 07:57:11 -07:00
committed by mikereiche
parent 08c7a5c7b4
commit 4c428c5351
3 changed files with 26 additions and 3 deletions

View File

@@ -16,7 +16,7 @@
package org.springframework.data.couchbase.repository.support;
import static org.springframework.data.couchbase.repository.support.Util.*;
import static org.springframework.data.couchbase.repository.support.Util.hasNonZeroVersionProperty;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@@ -26,7 +26,6 @@ import java.util.Objects;
import java.util.stream.Collectors;
import org.reactivestreams.Publisher;
import org.springframework.data.couchbase.core.CouchbaseOperations;
import org.springframework.data.couchbase.core.ReactiveCouchbaseOperations;
import org.springframework.data.couchbase.core.query.Query;
@@ -189,7 +188,8 @@ public class SimpleReactiveCouchbaseRepository<T, ID> implements ReactiveCouchba
@Override
public Mono<Void> deleteAll() {
return operations.removeByQuery(entityInformation.getJavaType()).all().then();
return operations.removeByQuery(entityInformation.getJavaType()).withConsistency(buildQueryScanConsistency()).all()
.then();
}
/**

View File

@@ -45,6 +45,10 @@ public interface ReactiveAirportRepository extends ReactiveSortingRepository<Air
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Flux<Airport> findAll();
@Override
@ScanConsistency(query = QueryScanConsistency.REQUEST_PLUS)
Mono<Void> deleteAll();
@Override
Mono<Airport> save(Airport a);

View File

@@ -195,6 +195,25 @@ public class ReactiveCouchbaseRepositoryQueryIntegrationTests extends JavaIntegr
}
}
@Test
void deleteAll() {
Airport vienna = new Airport("airports::vie", "vie", "LOWW");
Airport frankfurt = new Airport("airports::fra", "fra", "EDDF");
Airport losAngeles = new Airport("airports::lax", "lax", "KLAX");
try {
airportRepository.saveAll(asList(vienna, frankfurt, losAngeles)).as(StepVerifier::create)
.expectNext(vienna, frankfurt, losAngeles).verifyComplete();
airportRepository.deleteAll().as(StepVerifier::create).verifyComplete();
airportRepository.findAll().as(StepVerifier::create).verifyComplete();
} finally {
airportRepository.deleteAll().block();
}
}
@Configuration
@EnableReactiveCouchbaseRepositories("org.springframework.data.couchbase")
static class Config extends AbstractCouchbaseConfiguration {