diff --git a/spring-cql/src/main/java/org/springframework/cassandra/core/ReactiveCqlTemplate.java b/spring-cql/src/main/java/org/springframework/cassandra/core/ReactiveCqlTemplate.java index dd1db75c3..0bae2a3fb 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/core/ReactiveCqlTemplate.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/core/ReactiveCqlTemplate.java @@ -207,7 +207,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re Assert.notNull(action, "Callback object must not be null"); - return createFlux(action).onErrorResumeWith(translateException("ReactiveSessionCallback", getCql(action))); + return createFlux(action).onErrorMap(translateException("ReactiveSessionCallback", getCql(action))); } // ------------------------------------------------------------------------- @@ -241,7 +241,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re } return session.execute(stmt).flatMapMany(resultSetExtractor::extractData); - }).onErrorResumeWith(translateException("Query", cql)); + }).onErrorMap(translateException("Query", cql)); } /* (non-Javadoc) @@ -308,7 +308,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re } return session.execute(statement); - }).otherwise(translateException("QueryForResultSet", cql)); + }).onErrorMap(translateException("QueryForResultSet", cql)); } /* (non-Javadoc) @@ -317,7 +317,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re @Override public Flux queryForRows(String cql) throws DataAccessException { return queryForResultSet(cql).flatMapMany(ReactiveResultSet::rows) - .onErrorResumeWith(translateException("QueryForRows", cql)); + .onErrorMap(translateException("QueryForRows", cql)); } /* (non-Javadoc) @@ -362,7 +362,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re } return session.execute(stmt).flatMapMany(rse::extractData); - }).onErrorResumeWith(translateException("Query", statement.toString())); + }).onErrorMap(translateException("Query", statement.toString())); } /* (non-Javadoc) @@ -430,13 +430,13 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re } return session.execute(executedStatement); - }).otherwise(translateException("QueryForResultSet", statement.toString())); + }).onErrorMap(translateException("QueryForResultSet", statement.toString())); } @Override public Flux queryForRows(Statement statement) throws DataAccessException { return queryForResultSet(statement).flatMapMany(ReactiveResultSet::rows) - .onErrorResumeWith(translateException("QueryForRows", statement.toString())); + .onErrorMap(translateException("QueryForRows", statement.toString())); } // ------------------------------------------------------------------------- @@ -459,7 +459,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return psc.createPreparedStatement(session).doOnNext(this::applyStatementSettings) .flatMapMany(ps -> action.doInPreparedStatement(session, ps)); - }).onErrorResumeWith(translateException("ReactivePreparedStatementCallback", getCql(psc))); + }).onErrorMap(translateException("ReactivePreparedStatementCallback", getCql(psc))); } /* (non-Javadoc) @@ -498,7 +498,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re applyStatementSettings(boundStatement); return session.execute(boundStatement); - }).flatMap(rse::extractData)).onErrorResumeWith(translateException("Query", getCql(psc))); + }).flatMap(rse::extractData)).onErrorMap(translateException("Query", getCql(psc))); } /* (non-Javadoc) @@ -622,7 +622,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re @Override public Flux queryForRows(String cql, Object... args) throws DataAccessException { return queryForResultSet(cql, args).flatMapMany(ReactiveResultSet::rows) - .onErrorResumeWith(translateException("QueryForRows", cql)); + .onErrorMap(translateException("QueryForRows", cql)); } /* (non-Javadoc) @@ -728,17 +728,6 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re return Flux.defer(() -> callback.doInSession(session)); } - /** - * Exception translation {@link Function} intended for {@link Mono#otherwise(Function)} usage. - * - * @return the exception translation {@link Function} - */ - protected Function> translateException() { - - return throwable -> Mono.error( - throwable instanceof DriverException ? translateExceptionIfPossible((DriverException) throwable) : throwable); - } - /** * Exception translation {@link Function} intended for {@link Mono#otherwise(Function)} usage. * @@ -747,10 +736,9 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re * @return the exception translation {@link Function} * @see CqlProvider */ - protected Function> translateException(String task, String cql) { - - return throwable -> Mono - .error(throwable instanceof DriverException ? translate(task, cql, (DriverException) throwable) : throwable); + protected Function translateException(String task, String cql) { + return throwable -> throwable instanceof DriverException ? translate(task, cql, (DriverException) throwable) + : throwable; } /**