diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveSessionFactory.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveSessionFactory.java
index b889f7943..021b1315f 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveSessionFactory.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveSessionFactory.java
@@ -15,6 +15,8 @@
*/
package org.springframework.data.cassandra;
+import reactor.core.publisher.Mono;
+
/**
* Strategy interface to produce {@link ReactiveSession} instances.
*
@@ -36,5 +38,5 @@ public interface ReactiveSessionFactory {
*
* @return a {@link ReactiveSession}.
*/
- ReactiveSession getSession();
+ Mono getSession();
}
diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplate.java
index 6f3bc1837..76a81c3cd 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplate.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCqlTemplate.java
@@ -742,9 +742,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re
applyStatementSettings(statement);
- ReactiveSession session = getSession();
-
- return Flux.defer(() -> callback.doInStatement(session, statement));
+ return getSession().flatMapMany(session -> callback.doInStatement(session, statement));
}
/**
@@ -759,9 +757,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re
applyStatementSettings(statement);
- ReactiveSession session = getSession();
-
- return Mono.defer(() -> Mono.from(callback.doInStatement(session, statement)));
+ return getSession().flatMap(session -> Mono.from(callback.doInStatement(session, statement)));
}
/**
@@ -774,9 +770,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re
Assert.notNull(callback, "ReactiveStatementCallback must not be null");
- ReactiveSession session = getSession();
-
- return Flux.defer(() -> callback.doInSession(session));
+ return getSession().flatMapMany(callback::doInSession);
}
/**
@@ -860,7 +854,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re
return new ArgumentPreparedStatementBinder(args);
}
- private ReactiveSession getSession() {
+ private Mono getSession() {
ReactiveSessionFactory sessionFactory = getSessionFactory();
diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultReactiveSessionFactory.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultReactiveSessionFactory.java
index 30b2603f2..94aacbc14 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultReactiveSessionFactory.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultReactiveSessionFactory.java
@@ -15,6 +15,8 @@
*/
package org.springframework.data.cassandra.core.cql.session;
+import reactor.core.publisher.Mono;
+
import org.springframework.data.cassandra.ReactiveSession;
import org.springframework.data.cassandra.ReactiveSessionFactory;
import org.springframework.data.cassandra.core.cql.ReactiveRowMapperResultSetExtractor;
@@ -29,7 +31,7 @@ import org.springframework.data.cassandra.core.cql.ReactiveRowMapperResultSetExt
*/
public class DefaultReactiveSessionFactory implements ReactiveSessionFactory {
- private final ReactiveSession session;
+ private final Mono session;
/**
* Create a new {@link ReactiveRowMapperResultSetExtractor}.
@@ -37,11 +39,11 @@ public class DefaultReactiveSessionFactory implements ReactiveSessionFactory {
* @param session the {@link ReactiveSession} provides connections to Cassandra, must not be {@literal null}.
*/
public DefaultReactiveSessionFactory(ReactiveSession session) {
- this.session = session;
+ this.session = Mono.just(session);
}
@Override
- public ReactiveSession getSession() {
+ public Mono getSession() {
return session;
}
}
diff --git a/src/main/asciidoc/reference/migration-guide-2.2-to-3.0.adoc b/src/main/asciidoc/reference/migration-guide-2.2-to-3.0.adoc
index 93002bbac..099b52ed3 100644
--- a/src/main/asciidoc/reference/migration-guide-2.2-to-3.0.adoc
+++ b/src/main/asciidoc/reference/migration-guide-2.2-to-3.0.adoc
@@ -140,6 +140,8 @@ Paging state now uses `ByteBuffer`.
* Introduction of `StatementBuilder` to functionally build statements as the QueryBuilder API uses immutable statement types.
* `Session` bean renamed from `session` to `cassandraSession` and `SessionFactory` bean renamed from `sessionFactory` to `cassandraSessionFactory`.
* `ReactiveSession` bean renamed from `reactiveSession` to `reactiveCassandraSession` and `ReactiveSessionFactory` bean renamed from `reactiveSessionFactory` to `reactiveCassandraSessionFactory`.
+* `ReactiveSessionFactory.getSession()` now returns a `Mono`.
+Previously it returned just `ReactiveSession`.
* Data type resolution was moved into `ColumnTypeResolver` so all `DataType`-related methods were moved from `CassandraPersistentEntity`/`CassandraPersistentProperty` into `ColumnTypeResolver` (affected methods are `MappingContext.getDataType(…)`, `CassandraPersistentProperty.getDataType()`, `CassandraPersistentEntity.getUserType()`, and `CassandraPersistentEntity.getTupleType()`).
* Schema creation was moved from `MappingContext` to `SchemaFactory` (affected methods are `CassandraMappingContext.getCreateTableSpecificationFor(…)`, `CassandraMappingContext.getCreateIndexSpecificationsFor(…)`, and `CassandraMappingContext.getCreateUserTypeSpecificationFor(…)`).