From d9ddf0de1695f2d34e18eec2e4db6979727cc729 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Mon, 25 Oct 2021 09:10:02 +0200 Subject: [PATCH] Polishing. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Consistently use toCql(…) and String.format for CQL retrieval. See #1186 --- .../cassandra/core/cql/AsyncCqlTemplate.java | 30 ++++----- .../data/cassandra/core/cql/CqlTemplate.java | 24 +++---- .../core/cql/ReactiveCassandraAccessor.java | 16 ++++- .../core/cql/ReactiveCqlTemplate.java | 66 ++++++++----------- .../DefaultBridgedReactiveSession.java | 26 ++++---- 5 files changed, 80 insertions(+), 82 deletions(-) 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 0818548c2..c608dc9bf 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,6 +22,14 @@ 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; @@ -32,14 +40,6 @@ 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 @@ -163,7 +163,7 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera try { if (logger.isDebugEnabled()) { - logger.debug("Executing CQL statement [{}]", cql); + logger.debug(String.format("Executing CQL statement [%s]", cql)); } CompletionStage results = getCurrentSession().executeAsync(applyStatementSettings(newStatement(cql))) @@ -285,7 +285,7 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera try { if (logger.isDebugEnabled()) { - logger.debug("Executing statement [{}]", QueryExtractorDelegate.getCql(statement)); + logger.debug(String.format("Executing statement [%s]", toCql(statement))); } CompletionStage results = getCurrentSession() // @@ -294,9 +294,9 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera .thenCompose(ListenableFuture::completable); return new CassandraFutureAdapter<>(results, - ex -> translateExceptionIfPossible("Query", statement.toString(), ex)); + ex -> translateExceptionIfPossible("Query", toCql(statement), ex)); } catch (DriverException e) { - throw translateException("Query", statement.toString(), e); + throw translateException("Query", toCql(statement), e); } } @@ -446,7 +446,7 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera toCql(preparedStatementCreator), ex); try { if (logger.isDebugEnabled()) { - logger.debug("Preparing statement [{}] using {}", toCql(preparedStatementCreator), preparedStatementCreator); + logger.debug(String.format("Preparing statement [%s] using %s", toCql(preparedStatementCreator), preparedStatementCreator)); } CqlSession currentSession = getCurrentSession(); @@ -517,7 +517,7 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera try { if (logger.isDebugEnabled()) { - logger.debug("Preparing statement [{}] using {}", toCql(preparedStatementCreator), preparedStatementCreator); + logger.debug(String.format("Preparing statement [%s] using %s", toCql(preparedStatementCreator), preparedStatementCreator)); } CqlSession session = getCurrentSession(); @@ -525,7 +525,7 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera ListenableFuture> statementFuture = new MappingListenableFutureAdapter<>( preparedStatementCreator.createPreparedStatement(session), preparedStatement -> { if (logger.isDebugEnabled()) { - logger.debug("Executing prepared statement [{}]", QueryExtractorDelegate.getCql(preparedStatement)); + logger.debug(String.format("Executing prepared statement [%s]", toCql(preparedStatement))); } return applyStatementSettings(psb != null ? psb.bindValues(preparedStatement) : preparedStatement.bind()); 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 e97f6489b..6b53b4e9b 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,12 +25,6 @@ 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; @@ -40,6 +34,12 @@ 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 @@ -166,7 +166,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { try { if (logger.isDebugEnabled()) { - logger.debug("Executing CQL statement [{}]", cql); + logger.debug(String.format("Executing CQL statement [%s]", cql)); } Statement statement = applyStatementSettings(newStatement(cql)); @@ -292,12 +292,12 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { try { if (logger.isDebugEnabled()) { - logger.debug("Executing statement [{}]", QueryExtractorDelegate.getCql(statement)); + logger.debug(String.format("Executing statement [%s]", toCql(statement))); } return resultSetExtractor.extractData(getCurrentSession().execute(applyStatementSettings(statement))); } catch (DriverException e) { - throw translateException("Query", statement.toString(), e); + throw translateException("Query", toCql(statement), e); } } @@ -453,7 +453,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { try { if (logger.isDebugEnabled()) { - logger.debug("Preparing statement [{}] using {}", toCql(preparedStatementCreator), preparedStatementCreator); + logger.debug(String.format("Preparing statement [%s] using %s", toCql(preparedStatementCreator), preparedStatementCreator)); } CqlSession session = getCurrentSession(); @@ -523,7 +523,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { try { if (logger.isDebugEnabled()) { - logger.debug("Preparing statement [{}] using {}", toCql(preparedStatementCreator), preparedStatementCreator); + logger.debug(String.format("Preparing statement [%s] using %s", toCql(preparedStatementCreator), preparedStatementCreator)); } CqlSession session = getCurrentSession(); @@ -531,7 +531,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { PreparedStatement preparedStatement = preparedStatementCreator.createPreparedStatement(session); if (logger.isDebugEnabled()) { - logger.debug("Executing prepared statement [{}]", QueryExtractorDelegate.getCql(preparedStatement)); + logger.debug(String.format("Executing prepared statement [%s]", QueryExtractorDelegate.getCql(preparedStatement))); } Statement boundStatement = applyStatementSettings( diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCassandraAccessor.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCassandraAccessor.java index 4cfa24286..eb2fb904a 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCassandraAccessor.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/ReactiveCassandraAccessor.java @@ -15,6 +15,7 @@ */ package org.springframework.data.cassandra.core.cql; +import com.datastax.oss.driver.api.core.DriverException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -25,8 +26,6 @@ import org.springframework.data.cassandra.ReactiveSessionFactory; import org.springframework.lang.Nullable; import org.springframework.util.Assert; -import com.datastax.oss.driver.api.core.DriverException; - /** * Base class for {@link ReactiveCqlTemplate} and other CQL-accessing DAO helpers, defining common properties such as * {@link ReactiveSessionFactory} and exception translator. @@ -149,4 +148,17 @@ public abstract class ReactiveCassandraAccessor implements InitializingBean { return getExceptionTranslator().translate(task, cql, ex); } + + /** + * Determine CQL from potential provider object. + * + * @param cqlProvider object that's potentially a {@link CqlProvider} + * @return the CQL string, or {@literal null} + * @see CqlProvider + * @since 3.2.7 + */ + @Nullable + protected static String toCql(@Nullable Object cqlProvider) { + return QueryExtractorDelegate.getCql(cqlProvider); + } } 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 0f9d11a76..ae4efd916 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,23 +15,9 @@ */ 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; @@ -43,6 +29,18 @@ 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 @@ -302,7 +300,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re Assert.notNull(action, "Callback object must not be null"); - return createFlux(action).onErrorMap(translateException("ReactiveSessionCallback", getCql(action))); + return createFlux(action).onErrorMap(translateException("ReactiveSessionCallback", toCql(action))); } // ------------------------------------------------------------------------- @@ -439,11 +437,12 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return createFlux(statement, (session, stmt) -> { if (logger.isDebugEnabled()) { - logger.debug("Executing statement [{}]", QueryExtractorDelegate.getCql(statement)); + logger.debug(String.format("Executing statement [%s]", toCql(statement))); } - return session.execute(applyStatementSettings(statement)).flatMapMany(rse::extractData); - }).onErrorMap(translateException("Query", statement.toString())); + return session.execute(applyStatementSettings(statement)) + .flatMapMany(rse::extractData); + }).onErrorMap(translateException("Query", toCql(statement))); } /* (non-Javadoc) @@ -506,17 +505,17 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return createMono(statement, (session, executedStatement) -> { if (logger.isDebugEnabled()) { - logger.debug("Executing statement [{}]", QueryExtractorDelegate.getCql(statement)); + logger.debug(String.format("Executing statement [%s]", toCql(statement))); } return session.execute(applyStatementSettings(executedStatement)); - }).onErrorMap(translateException("QueryForResultSet", statement.toString())); + }).onErrorMap(translateException("QueryForResultSet", toCql(statement))); } @Override public Flux queryForRows(Statement statement) throws DataAccessException { return queryForResultSet(statement).flatMapMany(ReactiveResultSet::rows) - .onErrorMap(translateException("QueryForRows", statement.toString())); + .onErrorMap(translateException("QueryForRows", toCql(statement))); } // ------------------------------------------------------------------------- @@ -535,10 +534,11 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return createFlux(session -> { - logger.debug("Preparing statement [{}] using {}", getCql(psc), psc); + logger.debug(String.format("Preparing statement [%s] using %s", toCql(psc), psc)); - return psc.createPreparedStatement(session).flatMapMany(ps -> action.doInPreparedStatement(session, ps)); - }).onErrorMap(translateException("ReactivePreparedStatementCallback", getCql(psc))); + return psc.createPreparedStatement(session) + .flatMapMany(ps -> action.doInPreparedStatement(session, ps)); + }).onErrorMap(translateException("ReactivePreparedStatementCallback", toCql(psc))); } /* (non-Javadoc) @@ -569,7 +569,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return execute(psc, (session, preparedStatement) -> Mono.just(preparedStatement).flatMapMany(pps -> { if (logger.isDebugEnabled()) { - logger.debug("Executing prepared statement [{}]", QueryExtractorDelegate.getCql(preparedStatement)); + logger.debug(String.format("Executing prepared statement [%s]", toCql(preparedStatement))); } BoundStatement boundStatement = (preparedStatementBinder != null @@ -577,7 +577,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re : preparedStatement.bind()); return session.execute(applyStatementSettings(boundStatement)); - }).flatMap(rse::extractData)).onErrorMap(translateException("Query", getCql(psc))); + }).flatMap(rse::extractData)).onErrorMap(translateException("Query", toCql(psc))); } /* (non-Javadoc) @@ -737,7 +737,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return execute(newReactivePreparedStatementCreator(cql), (session, ps) -> Flux.from(args).flatMap(objects -> { if (logger.isDebugEnabled()) { - logger.debug("Executing prepared CQL statement [{}]", cql); + logger.debug(String.format("Executing prepared CQL statement [%s]", cql)); } BoundStatement boundStatement = newArgPreparedStatementBinder(objects).bindValues(ps); @@ -908,18 +908,6 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return sessionFactory.getSession(); } - /** - * Determine CQL from potential provider object. - * - * @param cqlProvider object that's potentially a {@link CqlProvider} - * @return the CQL string, or {@literal null} - * @see CqlProvider - */ - @Nullable - private static String getCql(@Nullable Object cqlProvider) { - return QueryExtractorDelegate.getCql(cqlProvider); - } - static class SimpleReactivePreparedStatementCreator implements ReactivePreparedStatementCreator, CqlProvider { private final SimpleStatement statement; 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 86494f551..75d3ff941 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,24 +15,12 @@ */ 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; @@ -47,6 +35,16 @@ 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} @@ -166,7 +164,7 @@ public class DefaultBridgedReactiveSession implements ReactiveSession { return Mono.fromCompletionStage(() -> { if (logger.isDebugEnabled()) { - logger.debug("Executing statement [{}]", getCql(statement)); + logger.debug(String.format("Executing statement [%s]", getCql(statement))); } return this.session.executeAsync(statement); @@ -195,7 +193,7 @@ public class DefaultBridgedReactiveSession implements ReactiveSession { return Mono.fromCompletionStage(() -> { if (logger.isDebugEnabled()) { - logger.debug("Preparing statement [{}]", getCql(statement)); + logger.debug(String.format("Preparing statement [%s]", getCql(statement))); } return this.session.prepareAsync(statement);