DATACASS-453 - Upgrade to Reactor 3.1 M2.

Adopt to API change from Publisher.subscribe() to Publisher.toProcessor(). Adopt to changed reactor-test groupId. Provide mocks for calls that allowed previously null Publishers.
This commit is contained in:
Mark Paluch
2017-06-01 10:04:33 +02:00
parent 0503d6362f
commit 86c510a8c5
3 changed files with 14 additions and 4 deletions

View File

@@ -90,7 +90,7 @@
</dependency>
<dependency>
<groupId>io.projectreactor.addons</groupId>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-test</artifactId>
<version>${reactor}</version>
<scope>test</scope>

View File

@@ -52,9 +52,9 @@ class ReactiveCassandraParameterAccessor extends CassandraParametersParameterAcc
}
if (ReactiveWrappers.isSingleValueType(value.getClass())) {
subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Mono.class).subscribe());
subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Mono.class).toProcessor());
} else {
subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Flux.class).collectList().subscribe());
subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Flux.class).collectList().toProcessor());
}
}
}

View File

@@ -31,9 +31,13 @@ import org.mockito.junit.MockitoJUnitRunner;
import org.springframework.data.cql.core.session.DefaultBridgedReactiveSession;
import com.datastax.driver.core.Cluster;
import com.datastax.driver.core.PreparedStatement;
import com.datastax.driver.core.RegularStatement;
import com.datastax.driver.core.ResultSetFuture;
import com.datastax.driver.core.Session;
import com.datastax.driver.core.SimpleStatement;
import com.datastax.driver.core.Statement;
import com.google.common.util.concurrent.ListenableFuture;
/**
* Unit tests for {@link DefaultBridgedReactiveSession}.
@@ -43,13 +47,19 @@ import com.datastax.driver.core.Statement;
@RunWith(MockitoJUnitRunner.class)
public class DefaultBridgedReactiveSessionUnitTests {
@Mock private Session sessionMock;
@Mock Session sessionMock;
@Mock ResultSetFuture future;
@Mock ListenableFuture<PreparedStatement> preparedStatementFuture;
private DefaultBridgedReactiveSession reactiveSession;
@Before
public void before() throws Exception {
reactiveSession = new DefaultBridgedReactiveSession(sessionMock, Schedulers.immediate());
when(sessionMock.executeAsync(any(Statement.class))).thenReturn(future);
when(sessionMock.prepareAsync(any(RegularStatement.class))).thenReturn(preparedStatementFuture);
}
@Test // DATACASS-335