Lazily obtain CQL for exception translation.

Closes #1236
This commit is contained in:
Mark Paluch
2022-03-03 09:15:19 +01:00
parent f8fe364c49
commit 83563239b2

View File

@@ -20,6 +20,7 @@ import reactor.core.publisher.Mono;
import java.util.Map;
import java.util.function.Function;
import java.util.function.Supplier;
import org.reactivestreams.Publisher;
@@ -29,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;
@@ -302,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)));
}
// -------------------------------------------------------------------------
@@ -436,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<String> 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)
@@ -503,20 +507,22 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re
Assert.notNull(statement, "CQL Statement must not be null");
Lazy<String> 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<Row> queryForRows(Statement<?> statement) throws DataAccessException {
return queryForResultSet(statement).flatMapMany(ReactiveResultSet::rows)
.onErrorMap(translateException("QueryForRows", toCql(statement)));
.onErrorMap(translateException("QueryForRows", () -> toCql(statement)));
}
// -------------------------------------------------------------------------
@@ -533,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<String> 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)
@@ -579,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)
@@ -819,7 +827,20 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re
* @see CqlProvider
*/
protected Function<Throwable, Throwable> 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.3.3
* @see CqlProvider
*/
protected Function<Throwable, Throwable> translateException(String task, Supplier<String> cql) {
return throwable -> throwable instanceof DriverException ? translate(task, cql.get(), (DriverException) throwable)
: throwable;
}