From 8438a477aee0bbe37ae17a18aecb04ae913ca6e6 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 587684d9f..07f6d0faf 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 540505cec..5a7ca6af4 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 031a15553..50ee3da24 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 a1998ae6a..7e9ec6d05 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 210b7bd9a..53484a5b5 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 495e55b8e..151fdd38b 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,10 +15,24 @@ */ package org.springframework.data.cassandra.core.cql; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + import java.util.Map; import java.util.Optional; 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; @@ -30,18 +44,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 @@ -441,8 +443,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))); } @@ -537,8 +538,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}