DATACASS-398 - Support RxJava 2 repositories.

Add RxJava 2 dependency. Add test to verify RxJava 2 interoperability.

Original Pull Request: #95
This commit is contained in:
Mark Paluch
2017-02-08 11:19:42 +01:00
committed by Christoph Strobl
parent 05d2a38848
commit 4f6548e932
2 changed files with 111 additions and 7 deletions

View File

@@ -92,6 +92,13 @@
<optional>true</optional>
</dependency>
<dependency>
<groupId>io.reactivex.rxjava2</groupId>
<artifactId>rxjava</artifactId>
<version>${rxjava2}</version>
<optional>true</optional>
</dependency>
<!-- CDI -->
<dependency>
<groupId>javax.enterprise</groupId>

View File

@@ -17,6 +17,10 @@ package org.springframework.data.cassandra.repository;
import static org.assertj.core.api.Assertions.*;
import io.reactivex.Flowable;
import io.reactivex.Maybe;
import io.reactivex.observers.TestObserver;
import org.springframework.data.cassandra.core.ReactiveCassandraTemplate;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
@@ -39,6 +43,7 @@ import org.springframework.data.cassandra.repository.config.EnableReactiveCassan
import org.springframework.data.cassandra.test.integration.support.IntegrationTestConfig;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import org.springframework.data.repository.reactive.RxJava1CrudRepository;
import org.springframework.data.repository.reactive.RxJava2CrudRepository;
import org.springframework.stereotype.Repository;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -68,9 +73,11 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
}
@Autowired Session session;
@Autowired MixedPersonRepostitory reactiveRepository;
@Autowired ReactiveCassandraTemplate template;
@Autowired MixedPersonRepository reactiveRepository;
@Autowired PersonRepostitory reactivePersonRepostitory;
@Autowired RxJava1PersonRepostitory rxJava1PersonRepostitory;
@Autowired RxJava2PersonRepostitory rxJava2PersonRepostitory;
Person dave, oliver, carter, boyd;
@@ -118,7 +125,8 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
}
@Test // DATACASS-335
public void simpleRxJavaMethodsShouldWork() {
public void simpleRxJava1MethodsShouldWork() {
rxJava1PersonRepostitory.exists(dave.getId()) //
.test() //
.awaitTerminalEvent() //
@@ -128,7 +136,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
}
@Test // DATACASS-335
public void existsWithSingleRxJavaIdMethodsShouldWork() {
public void existsWithSingleRxJava1IdMethodsShouldWork() {
rxJava1PersonRepostitory.exists(Single.just(dave.getId())) //
.test() //
@@ -139,7 +147,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
}
@Test // DATACASS-335
public void singleRxJavaQueryMethodShouldWork() {
public void singleRxJava1QueryMethodShouldWork() {
rxJava1PersonRepostitory.findManyByLastname(dave.getLastname()) //
.test() //
@@ -150,7 +158,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
}
@Test // DATACASS-335
public void singleProjectedRxJavaQueryMethodShouldWork() {
public void singleProjectedRxJava1QueryMethodShouldWork() {
List<ProjectedPerson> values = rxJava1PersonRepostitory.findProjectedByLastname(carter.getLastname()) //
.test() //
@@ -165,7 +173,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
}
@Test // DATACASS-335
public void observableRxJavaQueryMethodShouldWork() {
public void observableRxJava1QueryMethodShouldWork() {
rxJava1PersonRepostitory.findByLastname(boyd.getLastname()) //
.test() //
@@ -175,6 +183,83 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
.assertCompleted();
}
@Test // DATACASS-398
public void simpleRxJava2MethodsShouldWork() {
TestObserver<Boolean> testObserver = rxJava2PersonRepostitory.exists(dave.getId()).test();
testObserver.awaitTerminalEvent();
testObserver.assertComplete();
testObserver.assertNoErrors();
testObserver.assertValue(true);
}
@Test // DATACASS-398
public void existsWithSingleRxJava2IdMethodsShouldWork() {
TestObserver<Boolean> testObserver = rxJava2PersonRepostitory.exists(io.reactivex.Single.just(dave.getId())).test();
testObserver.awaitTerminalEvent();
testObserver.assertComplete();
testObserver.assertNoErrors();
testObserver.assertValue(true);
}
@Test // DATACASS-398
public void flowableRxJava2QueryMethodShouldWork() {
io.reactivex.subscribers.TestSubscriber<Person> testSubscriber = rxJava2PersonRepostitory
.findManyByLastname(dave.getLastname()).test();
testSubscriber.awaitTerminalEvent();
testSubscriber.assertComplete();
testSubscriber.assertNoErrors();
testSubscriber.assertValueCount(2);
}
@Test // DATACASS-398
public void singleProjectedRxJava2QueryMethodShouldWork() {
TestObserver<ProjectedPerson> testObserver = rxJava2PersonRepostitory
.findProjectedByLastname(Maybe.just(carter.getLastname())).test();
testObserver.awaitTerminalEvent();
testObserver.assertComplete();
testObserver.assertNoErrors();
testObserver.assertValue(actual -> {
assertThat(actual.getFirstname()).isEqualTo(carter.getFirstname());
return true;
});
}
@Test // DATACASS-398
public void observableProjectedRxJava2QueryMethodShouldWork() {
TestObserver<ProjectedPerson> testObserver = rxJava2PersonRepostitory
.findProjectedByLastname(Single.just(carter.getLastname())).test();
testObserver.awaitTerminalEvent();
testObserver.assertComplete();
testObserver.assertNoErrors();
testObserver.assertValue(actual -> {
assertThat(actual.getFirstname()).isEqualTo(carter.getFirstname());
return true;
});
}
@Test // DATACASS-398
public void maybeRxJava2QueryMethodShouldWork() {
TestObserver<Person> testObserver = rxJava2PersonRepostitory.findByLastname(boyd.getLastname()).test();
testObserver.awaitTerminalEvent();
testObserver.assertComplete();
testObserver.assertNoErrors();
testObserver.assertValue(boyd);
}
@Test // DATACASS-335
public void mixedRepositoryShouldWork() {
@@ -214,7 +299,19 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
}
@Repository
interface MixedPersonRepostitory extends ReactiveCassandraRepository<Person, String> {
interface RxJava2PersonRepostitory extends RxJava2CrudRepository<Person, String> {
Flowable<Person> findManyByLastname(String lastname);
Maybe<Person> findByLastname(String lastname);
io.reactivex.Single<ProjectedPerson> findProjectedByLastname(Maybe<String> lastname);
io.reactivex.Observable<ProjectedPerson> findProjectedByLastname(Single<String> lastname);
}
@Repository
interface MixedPersonRepository extends ReactiveCassandraRepository<Person, String> {
Single<Person> findByLastname(String lastname);