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 f58698704..3824a2cba 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 @@ -162,7 +162,7 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera try { if (logger.isDebugEnabled()) { - logger.debug("Executing CQL Statement [{}]", cql); + logger.debug("Executing CQL statement [{}]", cql); } CompletionStage results = getCurrentSession().executeAsync(applyStatementSettings(newStatement(cql))) @@ -284,7 +284,7 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera try { if (logger.isDebugEnabled()) { - logger.debug("Executing CQL Statement [{}]", statement); + logger.debug("Executing statement [{}]", QueryExtractorDelegate.getCql(statement)); } CompletionStage results = getCurrentSession() // @@ -524,7 +524,7 @@ public class AsyncCqlTemplate extends CassandraAccessor implements AsyncCqlOpera ListenableFuture> statementFuture = new MappingListenableFutureAdapter<>( preparedStatementCreator.createPreparedStatement(session), preparedStatement -> { if (logger.isDebugEnabled()) { - logger.debug("Executing prepared statement [{}]", preparedStatement); + logger.debug("Executing prepared statement [{}]", QueryExtractorDelegate.getCql(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/CassandraAccessor.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraAccessor.java index 843ed641f..a62b19926 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraAccessor.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraAccessor.java @@ -16,7 +16,6 @@ package org.springframework.data.cassandra.core.cql; import java.util.Map; -import java.util.Optional; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -420,11 +419,6 @@ public class CassandraAccessor implements InitializingBean { */ @Nullable protected static String toCql(@Nullable Object cqlProvider) { - - return Optional.ofNullable(cqlProvider) // - .filter(o -> o instanceof CqlProvider) // - .map(o -> (CqlProvider) o) // - .map(CqlProvider::getCql) // - .orElse(null); + return QueryExtractorDelegate.getCql(cqlProvider); } } 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 227f29652..f08b204cc 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 @@ -161,7 +161,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { try { if (logger.isDebugEnabled()) { - logger.debug("Executing CQL Statement [{}]", cql); + logger.debug("Executing CQL statement [{}]", cql); } Statement statement = applyStatementSettings(newStatement(cql)); @@ -287,7 +287,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { try { if (logger.isDebugEnabled()) { - logger.debug("Executing CQL Statement [{}]", statement); + logger.debug("Executing statement [{}]", QueryExtractorDelegate.getCql(statement)); } return resultSetExtractor.extractData(getCurrentSession().execute(applyStatementSettings(statement))); @@ -506,7 +506,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { PreparedStatement preparedStatement = preparedStatementCreator.createPreparedStatement(session); if (logger.isDebugEnabled()) { - logger.debug("Executing prepared statement [{}]", preparedStatement); + logger.debug("Executing prepared statement [{}]", QueryExtractorDelegate.getCql(preparedStatement)); } Statement boundStatement = applyStatementSettings( diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryExtractorDelegate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryExtractorDelegate.java new file mode 100644 index 000000000..d8b46a315 --- /dev/null +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryExtractorDelegate.java @@ -0,0 +1,79 @@ +/* + * Copyright 2020 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.cassandra.core.cql; + +import org.springframework.lang.Nullable; + +import com.datastax.oss.driver.api.core.cql.BatchStatement; +import com.datastax.oss.driver.api.core.cql.BatchableStatement; +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.SimpleStatement; +import com.datastax.oss.driver.api.core.cql.Statement; + +/** + * Utility to extract CQL queries from a {@link Statement}. + * + * @author Mark Paluch + * @since 3.0.5 + */ +public class QueryExtractorDelegate { + + /** + * Try to extract the {@link SimpleStatement#getQuery() CQL query} from a statement object. + * + * @param statement the statement object. + * @return the CQL query when {@code statement} is not {@code null}. + */ + @Nullable + public static String getCql(@Nullable Object statement) { + + if (statement == null) { + return null; + } + + if (statement instanceof CqlProvider) { + return ((CqlProvider) statement).getCql(); + } + + if (statement instanceof SimpleStatement) { + return ((SimpleStatement) statement).getQuery(); + } + + if (statement instanceof PreparedStatement) { + return ((PreparedStatement) statement).getQuery(); + } + + if (statement instanceof BoundStatement) { + return getCql(((BoundStatement) statement).getPreparedStatement()); + } + + if (statement instanceof BatchStatement) { + + StringBuilder builder = new StringBuilder(); + + for (BatchableStatement batchableStatement : ((BatchStatement) statement)) { + + String query = getCql(batchableStatement); + builder.append(query).append(query.endsWith(";") ? "" : ";"); + } + + return builder.toString(); + } + + return String.format("Unknown: %s", statement); + } +} 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 76a81c3cd..3c1390eb8 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 @@ -405,7 +405,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return createFlux(statement, (session, stmt) -> { if (logger.isDebugEnabled()) { - logger.debug("Executing CQL Statement [{}]", statement); + logger.debug("Executing statement [{}]", QueryExtractorDelegate.getCql(statement)); } return session.execute(applyStatementSettings(statement)).flatMapMany(rse::extractData); @@ -472,8 +472,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return createMono(statement, (session, executedStatement) -> { if (logger.isDebugEnabled()) { - logger.debug("Executing CQL [{}]", executedStatement); - + logger.debug("Executing statement [{}]", QueryExtractorDelegate.getCql(statement)); } return session.execute(applyStatementSettings(executedStatement)); @@ -533,14 +532,15 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re Assert.notNull(psc, "ReactivePreparedStatementCreator must not be null"); Assert.notNull(rse, "ReactiveResultSetExtractor object must not be null"); - return execute(psc, (session, ps) -> Mono.just(ps).flatMapMany(pps -> { + return execute(psc, (session, preparedStatement) -> Mono.just(preparedStatement).flatMapMany(pps -> { if (logger.isDebugEnabled()) { - logger.debug("Executing Prepared CQL Statement [{}]", ps.getQuery()); + logger.debug("Executing prepared statement [{}]", QueryExtractorDelegate.getCql(preparedStatement)); } - BoundStatement boundStatement = (preparedStatementBinder != null ? preparedStatementBinder.bindValues(ps) - : ps.bind()); + BoundStatement boundStatement = (preparedStatementBinder != null + ? preparedStatementBinder.bindValues(preparedStatement) + : preparedStatement.bind()); return session.execute(applyStatementSettings(boundStatement)); }).flatMap(rse::extractData)).onErrorMap(translateException("Query", getCql(psc))); @@ -703,7 +703,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("Executing prepared CQL statement [{}]", cql); } BoundStatement boundStatement = newArgPreparedStatementBinder(objects).bindValues(ps); @@ -872,12 +872,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re */ @Nullable private static String getCql(@Nullable Object cqlProvider) { - - return Optional.ofNullable(cqlProvider) // - .filter(o -> o instanceof CqlProvider) // - .map(o -> (CqlProvider) o) // - .map(CqlProvider::getCql) // - .orElse(null); + return QueryExtractorDelegate.getCql(cqlProvider); } static class SimpleReactivePreparedStatementCreator implements ReactivePreparedStatementCreator, CqlProvider { 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 ba07c2427..61b7a5f3f 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 @@ -35,6 +35,9 @@ import org.springframework.util.Assert; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.context.DriverContext; import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.BatchStatement; +import com.datastax.oss.driver.api.core.cql.BatchableStatement; +import com.datastax.oss.driver.api.core.cql.BoundStatement; import com.datastax.oss.driver.api.core.cql.ColumnDefinitions; import com.datastax.oss.driver.api.core.cql.ExecutionInfo; import com.datastax.oss.driver.api.core.cql.PreparedStatement; @@ -144,7 +147,7 @@ public class DefaultBridgedReactiveSession implements ReactiveSession { return Mono.fromCompletionStage(() -> { if (logger.isDebugEnabled()) { - logger.debug("Executing Statement [{}]", statement); + logger.debug("Executing statement [{}]", getCql(statement)); } return this.session.executeAsync(statement); @@ -173,13 +176,43 @@ public class DefaultBridgedReactiveSession implements ReactiveSession { return Mono.fromCompletionStage(() -> { if (logger.isDebugEnabled()) { - logger.debug("Preparing Statement [{}]", statement); + logger.debug("Preparing statement [{}]", getCql(statement)); } return this.session.prepareAsync(statement); }); } + private static String getCql(Object statement) { + + if (statement instanceof SimpleStatement) { + return ((SimpleStatement) statement).getQuery(); + } + + if (statement instanceof PreparedStatement) { + return ((PreparedStatement) statement).getQuery(); + } + + if (statement instanceof BoundStatement) { + return getCql(((BoundStatement) statement).getPreparedStatement()); + } + + if (statement instanceof BatchStatement) { + + StringBuilder builder = new StringBuilder(); + + for (BatchableStatement batchableStatement : ((BatchStatement) statement)) { + + String query = getCql(batchableStatement); + builder.append(query).append(query.endsWith(";") ? "" : ";"); + } + + return builder.toString(); + } + + return String.format("Unknown: %s", statement); + } + /* (non-Javadoc) * @see org.springframework.data.cassandra.ReactiveSession#close() */ @@ -293,4 +326,5 @@ public class DefaultBridgedReactiveSession implements ReactiveSession { return Collections.singletonList(getExecutionInfo()); } } + } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/QueryStatementCreator.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/QueryStatementCreator.java index 3ae39522c..50d8a6096 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/QueryStatementCreator.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/QueryStatementCreator.java @@ -24,6 +24,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.data.cassandra.core.StatementFactory; +import org.springframework.data.cassandra.core.cql.QueryExtractorDelegate; import org.springframework.data.cassandra.core.cql.QueryOptions; import org.springframework.data.cassandra.core.cql.QueryOptionsUtil; import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; @@ -110,7 +111,7 @@ class QueryStatementCreator { SimpleStatement statement = statementFactory.count(query, getPersistentEntity()).build(); if (LOG.isDebugEnabled()) { - LOG.debug(String.format("Created query [%s].", statement)); + LOG.debug(String.format("Created query [%s].", QueryExtractorDelegate.getCql(statement))); } return statement; @@ -137,7 +138,7 @@ class QueryStatementCreator { SimpleStatement statement = statementFactory.delete(query, getPersistentEntity()).build(); if (LOG.isDebugEnabled()) { - LOG.debug(String.format("Created query [%s].", statement)); + LOG.debug(String.format("Created query [%s].", QueryExtractorDelegate.getCql(statement))); } return statement; @@ -164,7 +165,7 @@ class QueryStatementCreator { SimpleStatement statement = statementFactory.select(query.limit(1), getPersistentEntity()).build(); if (LOG.isDebugEnabled()) { - LOG.debug(String.format("Created query [%s].", statement)); + LOG.debug(String.format("Created query [%s].", QueryExtractorDelegate.getCql(statement))); } return statement; @@ -248,7 +249,7 @@ class QueryStatementCreator { } if (LOG.isDebugEnabled()) { - LOG.debug(String.format("Created query [%s].", queryToUse)); + LOG.debug(String.format("Created query [%s].", QueryExtractorDelegate.getCql(queryToUse))); } return queryToUse;