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 24ce07d87..dd1db75c3 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 @@ -240,7 +240,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re logger.debug("Executing CQL Statement [{}]", cql); } - return session.execute(stmt).flatMap(resultSetExtractor::extractData); + return session.execute(stmt).flatMapMany(resultSetExtractor::extractData); }).onErrorResumeWith(translateException("Query", cql)); } @@ -316,7 +316,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re */ @Override public Flux queryForRows(String cql) throws DataAccessException { - return queryForResultSet(cql).flatMap(ReactiveResultSet::rows) + return queryForResultSet(cql).flatMapMany(ReactiveResultSet::rows) .onErrorResumeWith(translateException("QueryForRows", cql)); } @@ -361,7 +361,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re logger.debug("Executing CQL Statement [{}]", statement); } - return session.execute(stmt).flatMap(rse::extractData); + return session.execute(stmt).flatMapMany(rse::extractData); }).onErrorResumeWith(translateException("Query", statement.toString())); } @@ -435,7 +435,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re @Override public Flux queryForRows(Statement statement) throws DataAccessException { - return queryForResultSet(statement).flatMap(ReactiveResultSet::rows) + return queryForResultSet(statement).flatMapMany(ReactiveResultSet::rows) .onErrorResumeWith(translateException("QueryForRows", statement.toString())); } @@ -458,7 +458,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re logger.debug("Preparing statement [{}] using {}", getCql(psc), psc); return psc.createPreparedStatement(session).doOnNext(this::applyStatementSettings) - .flatMap(ps -> action.doInPreparedStatement(session, ps)); + .flatMapMany(ps -> action.doInPreparedStatement(session, ps)); }).onErrorResumeWith(translateException("ReactivePreparedStatementCallback", getCql(psc))); } @@ -486,7 +486,7 @@ 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).flatMap(pps -> { + return execute(psc, (session, ps) -> Mono.just(ps).flatMapMany(pps -> { if (logger.isDebugEnabled()) { logger.debug("Executing Prepared CQL Statement [{}]", ps.getQueryString()); @@ -621,7 +621,7 @@ public class ReactiveCqlTemplate extends ReactiveCassandraAccessor implements Re */ @Override public Flux queryForRows(String cql, Object... args) throws DataAccessException { - return queryForResultSet(cql, args).flatMap(ReactiveResultSet::rows) + return queryForResultSet(cql, args).flatMapMany(ReactiveResultSet::rows) .onErrorResumeWith(translateException("QueryForRows", cql)); } diff --git a/spring-cql/src/test/java/org/springframework/cassandra/core/ReactiveCqlTemplateUnitTests.java b/spring-cql/src/test/java/org/springframework/cassandra/core/ReactiveCqlTemplateUnitTests.java index 035f5fdc3..4084908b0 100644 --- a/spring-cql/src/test/java/org/springframework/cassandra/core/ReactiveCqlTemplateUnitTests.java +++ b/spring-cql/src/test/java/org/springframework/cassandra/core/ReactiveCqlTemplateUnitTests.java @@ -164,7 +164,7 @@ public class ReactiveCqlTemplateUnitTests { Mono mono = reactiveCqlTemplate.queryForResultSet("SELECT * from USERS"); - StepVerifier.create(mono.flatMap(ReactiveResultSet::rows)).expectNextCount(3).verifyComplete(); + StepVerifier.create(mono.flatMapMany(ReactiveResultSet::rows)).expectNextCount(3).verifyComplete(); verify(session).execute(any(Statement.class)); }); @@ -367,7 +367,7 @@ public class ReactiveCqlTemplateUnitTests { StepVerifier .create(reactiveCqlTemplate.queryForResultSet(new SimpleStatement("SELECT * from USERS")) - .flatMap(ReactiveResultSet::rows)) // + .flatMapMany(ReactiveResultSet::rows)) // .expectNextCount(3) // .verifyComplete(); @@ -535,7 +535,7 @@ public class ReactiveCqlTemplateUnitTests { Flux flux = reactiveCqlTemplate.execute("SELECT * from USERS", (session, ps) -> { - return session.execute(ps.bind("A")).flatMap(ReactiveResultSet::rows); + return session.execute(ps.bind("A")).flatMapMany(ReactiveResultSet::rows); }); StepVerifier.create(flux).expectNextCount(3).verifyComplete(); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java index a23ef7849..c227b83c9 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * Copyright 2016-2017 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. @@ -155,7 +155,7 @@ public class SimpleReactiveCassandraRepository Assert.notNull(mono, "The given id must not be null"); - return mono.then(id -> operations.selectOneById(id, entityInformation.getJavaType())); + return mono.flatMap(id -> operations.selectOneById(id, entityInformation.getJavaType())); } /* (non-Javadoc) @@ -177,7 +177,7 @@ public class SimpleReactiveCassandraRepository Assert.notNull(mono, "The given id must not be null"); - return mono.then(id -> operations.exists(id, entityInformation.getJavaType())); + return mono.flatMap(id -> operations.exists(id, entityInformation.getJavaType())); } /* (non-Javadoc)