diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraBatchTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraBatchTemplate.java index 2be9df9be..0eb579118 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraBatchTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraBatchTemplate.java @@ -15,17 +15,15 @@ */ package org.springframework.data.cassandra.core; -import reactor.core.publisher.Flux; -import reactor.core.publisher.Mono; - import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.function.Function; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; import org.springframework.data.cassandra.core.convert.CassandraConverter; import org.springframework.data.cassandra.core.convert.UpdateMapper; import org.springframework.data.cassandra.core.cql.WriteOptions; @@ -131,7 +129,7 @@ class ReactiveCassandraBatchTemplate implements ReactiveCassandraBatchOperations if (this.executed.compareAndSet(false, true)) { return Flux.merge(this.batchMonos) // - .flatMapIterable(Function.identity()) // + .flatMap(Flux::fromIterable) // .collectList() // .flatMap(statements -> { @@ -139,7 +137,8 @@ class ReactiveCassandraBatchTemplate implements ReactiveCassandraBatchOperations return this.operations.getReactiveCqlOperations().queryForResultSet(this.batch.build()); - }).flatMap(resultSet -> resultSet.rows().collectList() + }) // + .flatMap(resultSet -> resultSet.rows().collectList() .map(rows -> new WriteResult(resultSet.getAllExecutionInfo(), resultSet.wasApplied(), rows))); } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/convert/CassandraTypeMappingIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/convert/CassandraTypeMappingIntegrationTests.java index c04a8c5df..6622ecffa 100755 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/convert/CassandraTypeMappingIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/convert/CassandraTypeMappingIntegrationTests.java @@ -29,6 +29,7 @@ import java.nio.ByteBuffer; import java.time.Duration; import java.time.LocalDate; import java.time.LocalTime; +import java.time.temporal.ChronoUnit; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -616,7 +617,8 @@ public class CassandraTypeMappingIntegrationTests extends AbstractKeyspaceCreati AllPossibleTypes loaded = load(entity); - assertThat(loaded.getInstant()).isEqualTo(entity.getInstant()); + assertThat(loaded.getInstant().truncatedTo(ChronoUnit.MILLIS)) + .isEqualTo(entity.getInstant().truncatedTo(ChronoUnit.MILLIS)); } @Test // DATACASS-296 diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/mapping/event/CassandraTemplateEventIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/mapping/event/CassandraTemplateEventIntegrationTests.java index 5cc5d5a35..c44e3f5b8 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/mapping/event/CassandraTemplateEventIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/mapping/event/CassandraTemplateEventIntegrationTests.java @@ -18,6 +18,7 @@ package org.springframework.data.cassandra.core.mapping.event; import static org.assertj.core.api.Assertions.*; import java.util.List; +import java.util.stream.Collectors; import org.junit.Before; import org.junit.Test; @@ -50,7 +51,7 @@ public class CassandraTemplateEventIntegrationTests extends EventListenerIntegra @Test // DATACASS-106 public void streamShouldEmitEvents() { - template.stream("SELECT * FROM users;", User.class).count(); // Just load entire stream. + template.stream("SELECT * FROM users;", User.class).collect(Collectors.toList()); // Just load entire stream. assertThat(getListener().getAfterLoad()).extracting(CassandraMappingEvent::getTableName) .contains(CqlIdentifier.fromCql("users"));