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.
This commit is contained in:
committed by
Mark Paluch
parent
040a5938a3
commit
6f6d91c7d6
@@ -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)));
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"));
|
||||
|
||||
Reference in New Issue
Block a user