From 6f6d91c7d6b71a292df5a3ef03748ab2142ebe9b Mon Sep 17 00:00:00 2001 From: Greg Turnquist Date: Tue, 12 Nov 2019 11:40:53 -0600 Subject: [PATCH] DATACASS-699 - Address changes in Java 9 and newer. Adapt to changed Stream behavior that mapper is not used when issuing a count() operation. Fix generics. --- .../core/ReactiveCassandraBatchTemplate.java | 11 +++++------ .../convert/CassandraTypeMappingIntegrationTests.java | 4 +++- .../event/CassandraTemplateEventIntegrationTests.java | 3 ++- 3 files changed, 10 insertions(+), 8 deletions(-) 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"));