diff --git a/spring-data-cassandra/pom.xml b/spring-data-cassandra/pom.xml
index 7d3dd7e8a..a7209cefd 100644
--- a/spring-data-cassandra/pom.xml
+++ b/spring-data-cassandra/pom.xml
@@ -92,6 +92,13 @@
true
+
+ io.reactivex.rxjava2
+ rxjava
+ ${rxjava2}
+ true
+
+
javax.enterprise
diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ConvertingReactiveCassandraRepositoryTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ConvertingReactiveCassandraRepositoryTests.java
index 48a8399d4..7ab7c2537 100644
--- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ConvertingReactiveCassandraRepositoryTests.java
+++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ConvertingReactiveCassandraRepositoryTests.java
@@ -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 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 testObserver = rxJava2PersonRepostitory.exists(dave.getId()).test();
+
+ testObserver.awaitTerminalEvent();
+ testObserver.assertComplete();
+ testObserver.assertNoErrors();
+ testObserver.assertValue(true);
+ }
+
+ @Test // DATACASS-398
+ public void existsWithSingleRxJava2IdMethodsShouldWork() {
+
+ TestObserver 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 testSubscriber = rxJava2PersonRepostitory
+ .findManyByLastname(dave.getLastname()).test();
+
+ testSubscriber.awaitTerminalEvent();
+ testSubscriber.assertComplete();
+ testSubscriber.assertNoErrors();
+ testSubscriber.assertValueCount(2);
+ }
+
+ @Test // DATACASS-398
+ public void singleProjectedRxJava2QueryMethodShouldWork() {
+
+ TestObserver 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 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 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 {
+ interface RxJava2PersonRepostitory extends RxJava2CrudRepository {
+
+ Flowable findManyByLastname(String lastname);
+
+ Maybe findByLastname(String lastname);
+
+ io.reactivex.Single findProjectedByLastname(Maybe lastname);
+
+ io.reactivex.Observable findProjectedByLastname(Single lastname);
+ }
+
+ @Repository
+ interface MixedPersonRepository extends ReactiveCassandraRepository {
Single findByLastname(String lastname);