diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java index 6de0097cf..aa2284f7a 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java @@ -162,14 +162,14 @@ public class SimpleReactiveCassandraRepository implements ReactiveCassand } /* (non-Javadoc) - * @see org.springframework.data.repository.reactive.ReactiveCrudRepository#findById(reactor.core.publisher.Mono) + * @see org.springframework.data.repository.reactive.ReactiveCrudRepository#findById(org.reactivestreams.Publisher) */ @Override - public Mono findById(Mono mono) { + public Mono findById(Publisher publisher) { - Assert.notNull(mono, "The given id must not be null"); + Assert.notNull(publisher, "The given id must not be null"); - return mono.flatMap(id -> operations.selectOneById(id, entityInformation.getJavaType())); + return Mono.from(publisher).flatMap(id -> operations.selectOneById(id, entityInformation.getJavaType())); } /* (non-Javadoc) @@ -184,14 +184,14 @@ public class SimpleReactiveCassandraRepository implements ReactiveCassand } /* (non-Javadoc) - * @see org.springframework.data.repository.reactive.ReactiveCrudRepository#existsById(reactor.core.publisher.Mono) + * @see org.springframework.data.repository.reactive.ReactiveCrudRepository#existsById(org.reactivestreams.Publisher) */ @Override - public Mono existsById(Mono mono) { + public Mono existsById(Publisher publisher) { - Assert.notNull(mono, "The given id must not be null"); + Assert.notNull(publisher, "The given id must not be null"); - return mono.flatMap(id -> operations.exists(id, entityInformation.getJavaType())); + return Mono.from(publisher).flatMap(id -> operations.exists(id, entityInformation.getJavaType())); } /* (non-Javadoc) diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepositoryIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepositoryIntegrationTests.java index dc57ed882..d2d76f365 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepositoryIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepositoryIntegrationTests.java @@ -127,6 +127,15 @@ public class SimpleReactiveCassandraRepositoryIntegrationTests extends AbstractK StepVerifier.create(repository.existsById(Mono.just(dave.getId()))).expectNext(true).verifyComplete(); } + @Test // DATACASS-335 + public void existsByFluxOfIdShouldReturnTrueForExistingObject() { + + insertTestData(); + + StepVerifier.create(repository.existsById(Flux.just(dave.getId(), oliver.getId()))).expectNext(true) + .verifyComplete(); + } + @Test // DATACASS-335 public void existsByEmptyMonoOfIdShouldReturnEmptyMono() { StepVerifier.create(repository.existsById(Mono.empty())).verifyComplete(); @@ -153,6 +162,14 @@ public class SimpleReactiveCassandraRepositoryIntegrationTests extends AbstractK StepVerifier.create(repository.findById(Mono.just(dave.getId()))).expectNext(dave).verifyComplete(); } + @Test // DATACASS-462 + public void findByIdByFluxOfIdShouldReturnTrueForExistingObject() { + + insertTestData(); + + StepVerifier.create(repository.findById(Flux.just(dave.getId(), oliver.getId()))).expectNext(dave).verifyComplete(); + } + @Test // DATACASS-335 public void findByIdByEmptyMonoOfIdShouldReturnEmptyMono() { StepVerifier.create(repository.findById(Mono.empty())).verifyComplete();