diff --git a/spring-data-cassandra/pom.xml b/spring-data-cassandra/pom.xml
index 79a87991d..d37646895 100644
--- a/spring-data-cassandra/pom.xml
+++ b/spring-data-cassandra/pom.xml
@@ -86,23 +86,9 @@
- io.reactivex
+ io.reactivex.rxjava3
rxjava
- ${rxjava}
- true
-
-
-
- io.reactivex
- rxjava-reactive-streams
- ${rxjava-reactive-streams}
- true
-
-
-
- io.reactivex.rxjava2
- rxjava
- ${rxjava2}
+ ${rxjava3}
true
diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/ReactiveCassandraRepositoryFactoryBean.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/ReactiveCassandraRepositoryFactoryBean.java
index f2b7b9c0d..147cf6628 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/ReactiveCassandraRepositoryFactoryBean.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/ReactiveCassandraRepositoryFactoryBean.java
@@ -35,7 +35,6 @@ import org.springframework.util.Assert;
* @author Mark Paluch
* @since 2.0
* @see org.springframework.data.repository.reactive.ReactiveSortingRepository
- * @see org.springframework.data.repository.reactive.RxJava2SortingRepository
*/
public class ReactiveCassandraRepositoryFactoryBean, S, ID>
extends RepositoryFactoryBeanSupport {
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 da57034b7..5231187d3 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,13 +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.rxjava3.core.Single;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
-import rx.Observable;
-import rx.Single;
import java.util.Arrays;
import java.util.List;
@@ -40,7 +37,6 @@ import org.springframework.data.cassandra.repository.config.EnableReactiveCassan
import org.springframework.data.cassandra.repository.support.AbstractSpringDataEmbeddedCassandraIntegrationTest;
import org.springframework.data.cassandra.repository.support.IntegrationTestConfig;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
-import org.springframework.data.repository.reactive.RxJava2CrudRepository;
import org.springframework.stereotype.Repository;
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
@@ -71,8 +67,7 @@ class ConvertingReactiveCassandraRepositoryTests extends AbstractSpringDataEmbed
@Autowired CqlSession session;
@Autowired MixedUserRepository reactiveRepository;
@Autowired UserRepostitory reactiveUserRepostitory;
- @Autowired RxJava1UserRepository rxJava1UserRepository;
- @Autowired RxJava2UserRepository rxJava2UserRepository;
+ @Autowired RxJava3UserRepository rxJava3UserRepository;
private User dave;
private User oliver;
@@ -123,143 +118,72 @@ class ConvertingReactiveCassandraRepositoryTests extends AbstractSpringDataEmbed
}
@Test // DATACASS-335
- void simpleRxJava1MethodsShouldWork() {
+ void simpleRxJava3MethodsShouldWork() throws InterruptedException {
- rxJava1UserRepository.existsById(dave.getId()) //
+ rxJava3UserRepository.existsById(dave.getId()) //
.test() //
- .awaitTerminalEvent() //
+ .await() //
.assertResult(true) //
- .assertCompleted() //
+ .assertComplete() //
.assertNoErrors();
}
@Test // DATACASS-335
- void existsWithSingleRxJava1IdMethodsShouldWork() {
+ void existsWithSingleRxJava3IdMethodsShouldWork() throws InterruptedException {
- rxJava1UserRepository.existsById(Single.just(dave.getId())) //
+ rxJava3UserRepository.existsById(Single.just(dave.getId())) //
.test() //
- .awaitTerminalEvent() //
+ .await() //
.assertResult(true) //
- .assertCompleted() //
+ .assertComplete() //
.assertNoErrors();
}
@Test // DATACASS-335
- void singleRxJava1QueryMethodShouldWork() {
+ void singleRxJava3QueryMethodShouldWork() throws InterruptedException {
- rxJava1UserRepository.findManyByLastname(dave.getLastname()) //
+ rxJava3UserRepository.findManyByLastname(dave.getLastname()) //
.test() //
- .awaitTerminalEvent() //
+ .await() //
.assertValueCount(2) //
.assertNoErrors() //
- .assertCompleted();
+ .assertComplete();
}
@Test // DATACASS-335
- void singleProjectedRxJava1QueryMethodShouldWork() {
+ void singleProjectedRxJava3QueryMethodShouldWork() throws InterruptedException {
- List values = rxJava1UserRepository.findProjectedByLastname(carter.getLastname()) //
+ List values = rxJava3UserRepository.findProjectedByLastname(carter.getLastname()) //
.test() //
- .awaitTerminalEvent() //
+ .await() //
.assertValueCount(1) //
- .assertCompleted() //
+ .assertComplete() //
.assertNoErrors() //
- .getOnNextEvents();
+ .values();
ProjectedUser projectedUser = values.get(0);
assertThat(projectedUser.getFirstname()).isEqualTo(carter.getFirstname());
}
@Test // DATACASS-335
- void observableRxJava1QueryMethodShouldWork() {
+ void observableRxJava3QueryMethodShouldWork() throws InterruptedException {
- rxJava1UserRepository.findByLastname(boyd.getLastname()) //
+ rxJava3UserRepository.findByLastname(boyd.getLastname()) //
.test() //
- .awaitTerminalEvent() //
+ .await() //
.assertValue(boyd) //
.assertNoErrors() //
- .assertCompleted();
- }
-
- @Test // DATACASS-398
- void simpleRxJava2MethodsShouldWork() {
-
- rxJava2UserRepository.existsById(dave.getId()) //
- .test()//
- .assertValue(true) //
- .assertNoErrors() //
- .assertComplete() //
- .awaitTerminalEvent();
- }
-
- @Test // DATACASS-398
- void existsWithSingleRxJava2IdMethodsShouldWork() {
-
- rxJava2UserRepository.existsById(io.reactivex.Single.just(dave.getId())).test() //
- .assertValue(true) //
- .assertNoErrors() //
- .assertComplete() //
- .awaitTerminalEvent();
- }
-
- @Test // DATACASS-398
- void flowableRxJava2QueryMethodShouldWork() {
-
- rxJava2UserRepository.findManyByLastname(dave.getLastname()) //
- .test() //
- .assertValueCount(2) //
- .assertNoErrors() //
- .assertComplete() //
- .awaitTerminalEvent();
- }
-
- @Test // DATACASS-398
- void singleProjectedRxJava2QueryMethodShouldWork() {
-
- rxJava2UserRepository.findProjectedByLastname(Maybe.just(carter.getLastname())) //
- .test() //
- .assertValue(actual -> {
- assertThat(actual.getFirstname()).isEqualTo(carter.getFirstname());
- return true;
- }) //
- .assertComplete() //
- .assertNoErrors() //
- .awaitTerminalEvent();
- }
-
- @Test // DATACASS-398
- void observableProjectedRxJava2QueryMethodShouldWork() {
-
- rxJava2UserRepository.findProjectedByLastname(Single.just(carter.getLastname())) //
- .test() //
- .assertValue(actual -> {
- assertThat(actual.getFirstname()).isEqualTo(carter.getFirstname());
- return true;
- }) //
- .assertComplete() //
- .assertNoErrors() //
- .awaitTerminalEvent();
- }
-
- @Test // DATACASS-398
- void maybeRxJava2QueryMethodShouldWork() {
-
- rxJava2UserRepository.findByLastname(boyd.getLastname()) //
- .test() //
- .assertValue(boyd) //
- .assertNoErrors() //
- .assertComplete() //
- .awaitTerminalEvent();
+ .assertComplete();
}
@Test // DATACASS-335
- void mixedRepositoryShouldWork() {
+ void mixedRepositoryShouldWork() throws InterruptedException {
reactiveRepository.findByLastname(boyd.getLastname()) //
.test() //
- .awaitTerminalEvent() //
+ .await() //
.assertValue(boyd) //
- .assertCompleted() //
+ .assertComplete() //
.assertNoErrors();
}
@@ -280,29 +204,17 @@ class ConvertingReactiveCassandraRepositoryTests extends AbstractSpringDataEmbed
}
@Repository
- interface RxJava1UserRepository extends org.springframework.data.repository.Repository {
+ interface RxJava3UserRepository extends org.springframework.data.repository.Repository {
- Observable findManyByLastname(String lastname);
+ io.reactivex.rxjava3.core.Observable findManyByLastname(String lastname);
- Single findByLastname(String lastname);
+ io.reactivex.rxjava3.core.Single findByLastname(String lastname);
- Single findProjectedByLastname(String lastname);
+ io.reactivex.rxjava3.core.Single findProjectedByLastname(String lastname);
- Single existsById(String id);
+ io.reactivex.rxjava3.core.Single existsById(String id);
- Single existsById(Single id);
- }
-
- @Repository
- interface RxJava2UserRepository extends RxJava2CrudRepository {
-
- Flowable findManyByLastname(String lastname);
-
- Maybe findByLastname(String lastname);
-
- io.reactivex.Single findProjectedByLastname(Maybe lastname);
-
- io.reactivex.Observable findProjectedByLastname(Single lastname);
+ io.reactivex.rxjava3.core.Single existsById(io.reactivex.rxjava3.core.Single id);
}
@Repository
diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/config/ReactiveCassandraRepositoryConfigurationExtensionUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/config/ReactiveCassandraRepositoryConfigurationExtensionUnitTests.java
index 74b702add..92b6d50d8 100644
--- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/config/ReactiveCassandraRepositoryConfigurationExtensionUnitTests.java
+++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/config/ReactiveCassandraRepositoryConfigurationExtensionUnitTests.java
@@ -21,6 +21,7 @@ import java.util.Collection;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+
import org.springframework.beans.factory.support.BeanDefinitionRegistry;
import org.springframework.beans.factory.support.DefaultListableBeanFactory;
import org.springframework.core.env.Environment;
@@ -34,7 +35,7 @@ import org.springframework.data.repository.config.AnnotationRepositoryConfigurat
import org.springframework.data.repository.config.RepositoryConfiguration;
import org.springframework.data.repository.config.RepositoryConfigurationSource;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
-import org.springframework.data.repository.reactive.RxJava2CrudRepository;
+import org.springframework.data.repository.reactive.RxJava3CrudRepository;
/**
* Unit tests for {@link ReactiveCassandraRepositoryConfigurationExtension}.
@@ -107,7 +108,7 @@ public class ReactiveCassandraRepositoryConfigurationExtensionUnitTests {
@Table
private static class Sample {}
- interface SampleRepository extends RxJava2CrudRepository {}
+ interface SampleRepository extends RxJava3CrudRepository {}
interface UnannotatedRepository extends ReactiveCrudRepository