From 7f7ed92875290d139e97872ba90eb68cbf99b083 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Mon, 25 Oct 2021 09:53:39 +0200 Subject: [PATCH] Use nested classes for fetch size retrieval in session callbacks. We now use nested classes implementing CqlProvider to retrieve the fetch size in session callbacks. This enables CQL retrieval. Previously, CQL retrieval saw a class that didn't implement CqlProvider and reported therefore unknown CQL. See #1186 --- .../core/AsyncCassandraTemplate.java | 50 +++++++++------ .../cassandra/core/CassandraTemplate.java | 62 ++++++++++--------- .../core/ReactiveCassandraTemplate.java | 57 ++++++++++------- .../cassandra/core/cql/AsyncCqlTemplate.java | 25 ++++---- .../data/cassandra/core/cql/CqlTemplate.java | 12 ++-- .../core/cql/ReactiveCqlTemplate.java | 32 +++++----- .../DefaultBridgedReactiveSession.java | 22 ++++--- 7 files changed, 143 insertions(+), 117 deletions(-) diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java index 127fb0105..72d099247 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java @@ -22,6 +22,24 @@ import java.util.function.Function; import java.util.stream.Collectors; import java.util.stream.StreamSupport; +import com.datastax.oss.driver.api.core.CqlIdentifier; +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.DriverException; +import com.datastax.oss.driver.api.core.config.DefaultDriverOption; +import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.BoundStatement; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.cql.ResultSet; +import com.datastax.oss.driver.api.core.cql.Row; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.api.querybuilder.QueryBuilder; +import com.datastax.oss.driver.api.querybuilder.delete.Delete; +import com.datastax.oss.driver.api.querybuilder.insert.Insert; +import com.datastax.oss.driver.api.querybuilder.insert.RegularInsert; +import com.datastax.oss.driver.api.querybuilder.select.Select; +import com.datastax.oss.driver.api.querybuilder.truncate.Truncate; +import com.datastax.oss.driver.api.querybuilder.update.Update; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -63,25 +81,6 @@ import org.springframework.scheduling.annotation.AsyncResult; import org.springframework.util.Assert; import org.springframework.util.concurrent.ListenableFuture; -import com.datastax.oss.driver.api.core.CqlIdentifier; -import com.datastax.oss.driver.api.core.CqlSession; -import com.datastax.oss.driver.api.core.DriverException; -import com.datastax.oss.driver.api.core.config.DefaultDriverOption; -import com.datastax.oss.driver.api.core.cql.AsyncResultSet; -import com.datastax.oss.driver.api.core.cql.BoundStatement; -import com.datastax.oss.driver.api.core.cql.PreparedStatement; -import com.datastax.oss.driver.api.core.cql.ResultSet; -import com.datastax.oss.driver.api.core.cql.Row; -import com.datastax.oss.driver.api.core.cql.SimpleStatement; -import com.datastax.oss.driver.api.core.cql.Statement; -import com.datastax.oss.driver.api.querybuilder.QueryBuilder; -import com.datastax.oss.driver.api.querybuilder.delete.Delete; -import com.datastax.oss.driver.api.querybuilder.insert.Insert; -import com.datastax.oss.driver.api.querybuilder.insert.RegularInsert; -import com.datastax.oss.driver.api.querybuilder.select.Select; -import com.datastax.oss.driver.api.querybuilder.truncate.Truncate; -import com.datastax.oss.driver.api.querybuilder.update.Update; - /** * Primary implementation of {@link AsyncCassandraOperations}. It simplifies the use of asynchronous Cassandra usage and * helps to avoid common errors. It executes core Cassandra workflow. This class executes CQL queries or updates, @@ -933,9 +932,20 @@ public class AsyncCassandraTemplate return accessor.getFetchSize(); } } + class GetConfiguredPageSize implements AsyncSessionCallback, CqlProvider { + @Override + public ListenableFuture doInSession(CqlSession session) { + return AsyncResult.forValue(getConfiguredPageSize(session)); + } + + @Override + public String getCql() { + return QueryExtractorDelegate.getCql(statement); + } + } return getAsyncCqlOperations() - .execute((AsyncSessionCallback) session -> AsyncResult.forValue(getConfiguredPageSize(session))) + .execute(new GetConfiguredPageSize()) .completable().join(); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java index a576fc7da..6a3bbea98 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java @@ -20,6 +20,24 @@ import java.util.function.Consumer; import java.util.function.Function; import java.util.stream.Stream; +import com.datastax.oss.driver.api.core.CqlIdentifier; +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.DriverException; +import com.datastax.oss.driver.api.core.config.DefaultDriverOption; +import com.datastax.oss.driver.api.core.cql.BatchType; +import com.datastax.oss.driver.api.core.cql.BoundStatement; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.cql.ResultSet; +import com.datastax.oss.driver.api.core.cql.Row; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.api.querybuilder.QueryBuilder; +import com.datastax.oss.driver.api.querybuilder.delete.Delete; +import com.datastax.oss.driver.api.querybuilder.insert.Insert; +import com.datastax.oss.driver.api.querybuilder.insert.RegularInsert; +import com.datastax.oss.driver.api.querybuilder.select.Select; +import com.datastax.oss.driver.api.querybuilder.truncate.Truncate; +import com.datastax.oss.driver.api.querybuilder.update.Update; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -37,16 +55,7 @@ import org.springframework.data.cassandra.core.convert.CassandraConverter; import org.springframework.data.cassandra.core.convert.MappingCassandraConverter; import org.springframework.data.cassandra.core.convert.QueryMapper; import org.springframework.data.cassandra.core.convert.UpdateMapper; -import org.springframework.data.cassandra.core.cql.CassandraAccessor; -import org.springframework.data.cassandra.core.cql.CqlOperations; -import org.springframework.data.cassandra.core.cql.CqlProvider; -import org.springframework.data.cassandra.core.cql.CqlTemplate; -import org.springframework.data.cassandra.core.cql.PreparedStatementBinder; -import org.springframework.data.cassandra.core.cql.PreparedStatementCreator; -import org.springframework.data.cassandra.core.cql.QueryOptions; -import org.springframework.data.cassandra.core.cql.RowMapper; -import org.springframework.data.cassandra.core.cql.SingleColumnRowMapper; -import org.springframework.data.cassandra.core.cql.WriteOptions; +import org.springframework.data.cassandra.core.cql.*; import org.springframework.data.cassandra.core.cql.session.DefaultSessionFactory; import org.springframework.data.cassandra.core.cql.util.StatementBuilder; import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; @@ -69,25 +78,6 @@ import org.springframework.data.projection.SpelAwareProxyProjectionFactory; import org.springframework.lang.Nullable; import org.springframework.util.Assert; -import com.datastax.oss.driver.api.core.CqlIdentifier; -import com.datastax.oss.driver.api.core.CqlSession; -import com.datastax.oss.driver.api.core.DriverException; -import com.datastax.oss.driver.api.core.config.DefaultDriverOption; -import com.datastax.oss.driver.api.core.cql.BatchType; -import com.datastax.oss.driver.api.core.cql.BoundStatement; -import com.datastax.oss.driver.api.core.cql.PreparedStatement; -import com.datastax.oss.driver.api.core.cql.ResultSet; -import com.datastax.oss.driver.api.core.cql.Row; -import com.datastax.oss.driver.api.core.cql.SimpleStatement; -import com.datastax.oss.driver.api.core.cql.Statement; -import com.datastax.oss.driver.api.querybuilder.QueryBuilder; -import com.datastax.oss.driver.api.querybuilder.delete.Delete; -import com.datastax.oss.driver.api.querybuilder.insert.Insert; -import com.datastax.oss.driver.api.querybuilder.insert.RegularInsert; -import com.datastax.oss.driver.api.querybuilder.select.Select; -import com.datastax.oss.driver.api.querybuilder.truncate.Truncate; -import com.datastax.oss.driver.api.querybuilder.update.Update; - /** * Primary implementation of {@link CassandraOperations}. It simplifies the use of Cassandra usage and helps to avoid * common errors. It executes core Cassandra workflow. This class executes CQL queries or updates, initiating iteration @@ -992,7 +982,19 @@ public class CassandraTemplate implements CassandraOperations, ApplicationEventP } } - return getCqlOperations().execute(this::getConfiguredPageSize); + class GetConfiguredPageSize implements SessionCallback, CqlProvider { + @Override + public Integer doInSession(CqlSession session) { + return getConfiguredPageSize(session); + } + + @Override + public String getCql() { + return QueryExtractorDelegate.getCql(statement); + } + } + + return getCqlOperations().execute(new GetConfiguredPageSize()); } @SuppressWarnings("unchecked") diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java index 9bae0fab9..a9f1d0e2d 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java @@ -15,16 +15,33 @@ */ package org.springframework.data.cassandra.core; -import reactor.core.publisher.Flux; -import reactor.core.publisher.Mono; -import reactor.core.publisher.SynchronousSink; - import java.util.Collections; import java.util.function.BiConsumer; import java.util.function.Function; +import com.datastax.oss.driver.api.core.CqlIdentifier; +import com.datastax.oss.driver.api.core.DriverException; +import com.datastax.oss.driver.api.core.config.DefaultDriverOption; +import com.datastax.oss.driver.api.core.context.DriverContext; +import com.datastax.oss.driver.api.core.cql.BatchType; +import com.datastax.oss.driver.api.core.cql.BoundStatement; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.cql.Row; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.api.querybuilder.QueryBuilder; +import com.datastax.oss.driver.api.querybuilder.delete.Delete; +import com.datastax.oss.driver.api.querybuilder.insert.Insert; +import com.datastax.oss.driver.api.querybuilder.insert.RegularInsert; +import com.datastax.oss.driver.api.querybuilder.select.Select; +import com.datastax.oss.driver.api.querybuilder.truncate.Truncate; +import com.datastax.oss.driver.api.querybuilder.update.Update; +import org.reactivestreams.Publisher; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.core.publisher.SynchronousSink; import org.springframework.beans.BeansException; import org.springframework.context.ApplicationContext; @@ -65,24 +82,6 @@ import org.springframework.data.projection.SpelAwareProxyProjectionFactory; import org.springframework.lang.Nullable; import org.springframework.util.Assert; -import com.datastax.oss.driver.api.core.CqlIdentifier; -import com.datastax.oss.driver.api.core.DriverException; -import com.datastax.oss.driver.api.core.config.DefaultDriverOption; -import com.datastax.oss.driver.api.core.context.DriverContext; -import com.datastax.oss.driver.api.core.cql.BatchType; -import com.datastax.oss.driver.api.core.cql.BoundStatement; -import com.datastax.oss.driver.api.core.cql.PreparedStatement; -import com.datastax.oss.driver.api.core.cql.Row; -import com.datastax.oss.driver.api.core.cql.SimpleStatement; -import com.datastax.oss.driver.api.core.cql.Statement; -import com.datastax.oss.driver.api.querybuilder.QueryBuilder; -import com.datastax.oss.driver.api.querybuilder.delete.Delete; -import com.datastax.oss.driver.api.querybuilder.insert.Insert; -import com.datastax.oss.driver.api.querybuilder.insert.RegularInsert; -import com.datastax.oss.driver.api.querybuilder.select.Select; -import com.datastax.oss.driver.api.querybuilder.truncate.Truncate; -import com.datastax.oss.driver.api.querybuilder.update.Update; - /** * Primary implementation of {@link ReactiveCassandraOperations}. It simplifies the use of Reactive Cassandra usage and * helps to avoid common errors. It executes core Cassandra workflow. This class executes CQL queries or updates, @@ -937,8 +936,20 @@ public class ReactiveCassandraTemplate } } + class GetConfiguredPageSize implements ReactiveSessionCallback, CqlProvider { + @Override + public Publisher doInSession(ReactiveSession session) { + return Mono.just(getConfiguredPageSize(session.getContext())); + } + + @Override + public String getCql() { + return QueryExtractorDelegate.getCql(statement); + } + } + return getReactiveCqlOperations() - .execute((ReactiveSessionCallback) session -> Mono.just(getConfiguredPageSize(session.getContext()))) + .execute(new GetConfiguredPageSize()) .single(); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/AsyncCqlTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/AsyncCqlTemplate.java index c608dc9bf..f40ded36f 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/AsyncCqlTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/AsyncCqlTemplate.java @@ -22,14 +22,6 @@ import java.util.concurrent.CompletionStage; import java.util.concurrent.ExecutionException; import java.util.function.Function; -import com.datastax.oss.driver.api.core.CqlSession; -import com.datastax.oss.driver.api.core.DriverException; -import com.datastax.oss.driver.api.core.cql.AsyncResultSet; -import com.datastax.oss.driver.api.core.cql.PreparedStatement; -import com.datastax.oss.driver.api.core.cql.ResultSet; -import com.datastax.oss.driver.api.core.cql.SimpleStatement; -import com.datastax.oss.driver.api.core.cql.Statement; - import org.springframework.dao.DataAccessException; import org.springframework.dao.support.DataAccessUtils; import org.springframework.dao.support.PersistenceExceptionTranslator; @@ -40,6 +32,14 @@ import org.springframework.util.Assert; import org.springframework.util.concurrent.ListenableFuture; import org.springframework.util.concurrent.SettableListenableFuture; +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.DriverException; +import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.cql.ResultSet; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import com.datastax.oss.driver.api.core.cql.Statement; + /** * This is the central class in the CQL core package for asynchronous Cassandra data access. It simplifies the * use of CQL and helps to avoid common errors. It executes core CQL workflow, leaving application code to provide CQL @@ -293,8 +293,7 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera .thenApply(resultSetExtractor::extractData) // .thenCompose(ListenableFuture::completable); - return new CassandraFutureAdapter<>(results, - ex -> translateExceptionIfPossible("Query", toCql(statement), ex)); + return new CassandraFutureAdapter<>(results, ex -> translateExceptionIfPossible("Query", toCql(statement), ex)); } catch (DriverException e) { throw translateException("Query", toCql(statement), e); } @@ -446,7 +445,8 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera toCql(preparedStatementCreator), ex); try { if (logger.isDebugEnabled()) { - logger.debug(String.format("Preparing statement [%s] using %s", toCql(preparedStatementCreator), preparedStatementCreator)); + logger.debug(String.format("Preparing statement [%s] using %s", toCql(preparedStatementCreator), + preparedStatementCreator)); } CqlSession currentSession = getCurrentSession(); @@ -517,7 +517,8 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera try { if (logger.isDebugEnabled()) { - logger.debug(String.format("Preparing statement [%s] using %s", toCql(preparedStatementCreator), preparedStatementCreator)); + logger.debug(String.format("Preparing statement [%s] using %s", toCql(preparedStatementCreator), + preparedStatementCreator)); } CqlSession session = getCurrentSession(); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CqlTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CqlTemplate.java index 6b53b4e9b..a265fe49e 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CqlTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CqlTemplate.java @@ -25,6 +25,12 @@ import java.util.function.Function; import java.util.stream.Stream; import java.util.stream.StreamSupport; +import org.springframework.dao.DataAccessException; +import org.springframework.dao.support.DataAccessUtils; +import org.springframework.data.cassandra.SessionFactory; +import org.springframework.lang.Nullable; +import org.springframework.util.Assert; + import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.DriverException; import com.datastax.oss.driver.api.core.cql.PreparedStatement; @@ -34,12 +40,6 @@ import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.core.metadata.Node; -import org.springframework.dao.DataAccessException; -import org.springframework.dao.support.DataAccessUtils; -import org.springframework.data.cassandra.SessionFactory; -import org.springframework.lang.Nullable; -import org.springframework.util.Assert; - /** * This is the central class in the CQL core package. It simplifies the use of CQL and helps to avoid common * errors. It executes core CQL workflow, leaving application code to provide CQL and extract results. This class 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 ae4efd916..66a82ec76 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 @@ -15,9 +15,23 @@ */ package org.springframework.data.cassandra.core.cql; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + import java.util.Map; import java.util.function.Function; +import org.reactivestreams.Publisher; + +import org.springframework.dao.DataAccessException; +import org.springframework.dao.support.DataAccessUtils; +import org.springframework.data.cassandra.ReactiveResultSet; +import org.springframework.data.cassandra.ReactiveSession; +import org.springframework.data.cassandra.ReactiveSessionFactory; +import org.springframework.data.cassandra.core.cql.session.DefaultReactiveSessionFactory; +import org.springframework.lang.Nullable; +import org.springframework.util.Assert; + import com.datastax.oss.driver.api.core.ConsistencyLevel; import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.CqlSession; @@ -29,18 +43,6 @@ import com.datastax.oss.driver.api.core.cql.Row; import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.core.retry.RetryPolicy; -import org.reactivestreams.Publisher; -import reactor.core.publisher.Flux; -import reactor.core.publisher.Mono; - -import org.springframework.dao.DataAccessException; -import org.springframework.dao.support.DataAccessUtils; -import org.springframework.data.cassandra.ReactiveResultSet; -import org.springframework.data.cassandra.ReactiveSession; -import org.springframework.data.cassandra.ReactiveSessionFactory; -import org.springframework.data.cassandra.core.cql.session.DefaultReactiveSessionFactory; -import org.springframework.lang.Nullable; -import org.springframework.util.Assert; /** * This is the central class in the CQL core package for reactive Cassandra data access. It simplifies the use of @@ -440,8 +442,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re logger.debug(String.format("Executing statement [%s]", toCql(statement))); } - return session.execute(applyStatementSettings(statement)) - .flatMapMany(rse::extractData); + return session.execute(applyStatementSettings(statement)).flatMapMany(rse::extractData); }).onErrorMap(translateException("Query", toCql(statement))); } @@ -536,8 +537,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re logger.debug(String.format("Preparing statement [%s] using %s", toCql(psc), psc)); - return psc.createPreparedStatement(session) - .flatMapMany(ps -> action.doInPreparedStatement(session, ps)); + return psc.createPreparedStatement(session).flatMapMany(ps -> action.doInPreparedStatement(session, ps)); }).onErrorMap(translateException("ReactivePreparedStatementCallback", toCql(psc))); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java index 75d3ff941..cacf342bc 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java @@ -15,12 +15,24 @@ */ package org.springframework.data.cassandra.core.cql.session; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.core.publisher.MonoProcessor; +import reactor.core.scheduler.Scheduler; + import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.concurrent.CompletionStage; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.springframework.data.cassandra.ReactiveResultSet; +import org.springframework.data.cassandra.ReactiveSession; +import org.springframework.util.Assert; + import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.context.DriverContext; @@ -35,16 +47,6 @@ import com.datastax.oss.driver.api.core.cql.Row; import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.core.metadata.Metadata; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import reactor.core.publisher.Flux; -import reactor.core.publisher.Mono; -import reactor.core.publisher.MonoProcessor; -import reactor.core.scheduler.Scheduler; - -import org.springframework.data.cassandra.ReactiveResultSet; -import org.springframework.data.cassandra.ReactiveSession; -import org.springframework.util.Assert; /** * Default implementation of a {@link ReactiveSession}. This implementation bridges asynchronous {@link CqlSession}