From 1267514df28ee94997ad72a8a147e8878edd8b07 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Thu, 3 Mar 2022 09:15:19 +0100 Subject: [PATCH] Lazily obtain CQL for exception translation. Closes #1236 --- .../core/cql/ReactiveCqlTemplate.java | 42 ++++++++++++++----- 1 file changed, 31 insertions(+), 11 deletions(-) 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 3eff4b906..8d56574cc 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 @@ -19,8 +19,8 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.util.Map; -import java.util.Optional; import java.util.function.Function; +import java.util.function.Supplier; import org.reactivestreams.Publisher; @@ -30,6 +30,7 @@ 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.data.util.Lazy; import org.springframework.lang.Nullable; import org.springframework.util.Assert; @@ -303,7 +304,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re Assert.notNull(action, "Callback object must not be null"); - return createFlux(action).onErrorMap(translateException("ReactiveSessionCallback", toCql(action))); + return createFlux(action).onErrorMap(translateException("ReactiveSessionCallback", () -> toCql(action))); } // ------------------------------------------------------------------------- @@ -437,14 +438,16 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re Assert.notNull(statement, "CQL Statement must not be null"); Assert.notNull(rse, "ReactiveResultSetExtractor must not be null"); + Lazy cql = Lazy.of(() -> toCql(statement)); + return createFlux(statement, (session, stmt) -> { if (logger.isDebugEnabled()) { - logger.debug(String.format("Executing statement [%s]", toCql(statement))); + logger.debug(String.format("Executing statement [%s]", cql.get())); } return session.execute(applyStatementSettings(statement)).flatMapMany(rse::extractData); - }).onErrorMap(translateException("Query", toCql(statement))); + }).onErrorMap(translateException("Query", cql)); } /* (non-Javadoc) @@ -504,20 +507,22 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re Assert.notNull(statement, "CQL Statement must not be null"); + Lazy cql = Lazy.of(() -> toCql(statement)); + return createMono(statement, (session, executedStatement) -> { if (logger.isDebugEnabled()) { - logger.debug(String.format("Executing statement [%s]", toCql(statement))); + logger.debug(String.format("Executing statement [%s]", cql.get())); } return session.execute(applyStatementSettings(executedStatement)); - }).onErrorMap(translateException("QueryForResultSet", toCql(statement))); + }).onErrorMap(translateException("QueryForResultSet", cql)); } @Override public Flux queryForRows(Statement statement) throws DataAccessException { return queryForResultSet(statement).flatMapMany(ReactiveResultSet::rows) - .onErrorMap(translateException("QueryForRows", toCql(statement))); + .onErrorMap(translateException("QueryForRows", () -> toCql(statement))); } // ------------------------------------------------------------------------- @@ -534,14 +539,16 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re Assert.notNull(psc, "ReactivePreparedStatementCreator must not be null"); Assert.notNull(action, "ReactivePreparedStatementCallback object must not be null"); + Lazy cql = Lazy.of(() -> toCql(psc)); + return createFlux(session -> { if (logger.isDebugEnabled()) { - logger.debug(String.format("Preparing statement [%s] using %s", toCql(psc), psc)); + logger.debug(String.format("Preparing statement [%s] using %s", cql.get(), psc)); } return psc.createPreparedStatement(session).flatMapMany(ps -> action.doInPreparedStatement(session, ps)); - }).onErrorMap(translateException("ReactivePreparedStatementCallback", toCql(psc))); + }).onErrorMap(translateException("ReactivePreparedStatementCallback", cql)); } /* (non-Javadoc) @@ -580,7 +587,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re : preparedStatement.bind()); return session.execute(applyStatementSettings(boundStatement)); - }).flatMap(rse::extractData)).onErrorMap(translateException("Query", toCql(psc))); + }).flatMap(rse::extractData)).onErrorMap(translateException("Query", () -> toCql(psc))); } /* (non-Javadoc) @@ -820,7 +827,20 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re * @see CqlProvider */ protected Function translateException(String task, @Nullable String cql) { - return throwable -> throwable instanceof DriverException ? translate(task, cql, (DriverException) throwable) + return translateException(task, () -> cql); + } + + /** + * Exception translation {@link Function} intended for {@link Mono#onErrorMap(Function)} usage. + * + * @param task readable text describing the task being attempted + * @param cql supplier of CQL query or update that caused the problem (may be {@literal null}) + * @return the exception translation {@link Function} + * @since 3.2.10 + * @see CqlProvider + */ + protected Function translateException(String task, Supplier cql) { + return throwable -> throwable instanceof DriverException ? translate(task, cql.get(), (DriverException) throwable) : throwable; }