diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveRowMapperResultSetExtractor.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveRowMapperResultSetExtractor.java index d66de82d0..e60ed540a 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveRowMapperResultSetExtractor.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveRowMapperResultSetExtractor.java @@ -15,9 +15,8 @@ */ package org.springframework.data.cassandra.core.cql; -import reactor.core.publisher.Mono; - import org.reactivestreams.Publisher; + import org.springframework.dao.DataAccessException; import org.springframework.data.cassandra.ReactiveResultSet; import org.springframework.util.Assert; @@ -61,15 +60,13 @@ public class ReactiveRowMapperResultSetExtractor implements ReactiveResultSet @Override public Publisher extractData(ReactiveResultSet resultSet) throws DriverException, DataAccessException { - return resultSet.rows().flatMap(row -> { + return resultSet.rows().handle((row, sink) -> { T value = this.rowMapper.mapRow(row, 0); - if (value == null) { - return Mono.empty(); + if (value != null) { + sink.next(value); } - - return Mono.just(value); }); } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java index 878d24a05..fb5c63775 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java @@ -90,7 +90,7 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM users"); + assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users"); } @Test // DATACASS-335 @@ -124,7 +124,7 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); } @Test // DATACASS-696 @@ -144,7 +144,7 @@ public class ReactiveCassandraTemplateUnitTests { template.exists("myid", User.class).as(StepVerifier::create).expectNext(true).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); } @Test // DATACASS-335 @@ -155,7 +155,7 @@ public class ReactiveCassandraTemplateUnitTests { template.exists("myid", User.class).as(StepVerifier::create).expectNext(false).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); } @Test // DATACASS-512 @@ -166,7 +166,7 @@ public class ReactiveCassandraTemplateUnitTests { template.exists(Query.empty(), User.class).as(StepVerifier::create).expectNext(true).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM users LIMIT 1;"); + assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users LIMIT 1;"); } @Test // DATACASS-512 @@ -177,7 +177,7 @@ public class ReactiveCassandraTemplateUnitTests { template.exists(Query.empty(), User.class).as(StepVerifier::create).expectNext(false).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM users LIMIT 1;"); + assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users LIMIT 1;"); } @Test // DATACASS-335 @@ -190,7 +190,7 @@ public class ReactiveCassandraTemplateUnitTests { template.count(User.class).as(StepVerifier::create).expectNext(42L).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT count(*) FROM users;"); + assertThat(statementCaptor.getValue()).hasToString("SELECT count(*) FROM users;"); } @Test // DATACASS-512 @@ -203,7 +203,7 @@ public class ReactiveCassandraTemplateUnitTests { template.count(Query.empty(), User.class).as(StepVerifier::create).expectNext(42L).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT COUNT(1) FROM users;"); + assertThat(statementCaptor.getValue()).hasToString("SELECT COUNT(1) FROM users;"); } @Test // DATACASS-335 @@ -216,8 +216,8 @@ public class ReactiveCassandraTemplateUnitTests { template.insert(user).as(StepVerifier::create).expectNext(user).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()) - .isEqualTo("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White');"); + assertThat(statementCaptor.getValue()) + .hasToString("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White');"); } @Test // DATACASS-335 @@ -245,8 +245,8 @@ public class ReactiveCassandraTemplateUnitTests { template.update(user).as(StepVerifier::create).expectNext(user).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()) - .isEqualTo("UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue()) + .hasToString("UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg';"); } @Test // DATACASS-335 @@ -260,7 +260,7 @@ public class ReactiveCassandraTemplateUnitTests { template.delete(user).as(StepVerifier::create).expectNext(user).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("DELETE FROM users WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue()).hasToString("DELETE FROM users WHERE id='heisenberg';"); } @Test // DATACASS-335 @@ -269,6 +269,6 @@ public class ReactiveCassandraTemplateUnitTests { template.truncate(User.class).as(StepVerifier::create).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().toString()).isEqualTo("TRUNCATE users;"); + assertThat(statementCaptor.getValue()).hasToString("TRUNCATE users;"); } }