From da702c2a49abf1b0314ac2730702f43cfdc7976a Mon Sep 17 00:00:00 2001 From: Sebastien Deleuze Date: Wed, 4 Sep 2019 11:00:49 +0200 Subject: [PATCH] DATACASS-685 - Upgrade to Coroutines 1.3. Original pull request: #164. --- .../core/ReactiveSelectOperationExtensions.kt | 13 +++++-------- .../ReactiveSelectOperationExtensionsUnitTests.kt | 4 ++-- 2 files changed, 7 insertions(+), 10 deletions(-) diff --git a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensions.kt b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensions.kt index 43970e9ed..3e6ac79bf 100644 --- a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensions.kt +++ b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensions.kt @@ -15,11 +15,11 @@ */ package org.springframework.data.cassandra.core -import kotlinx.coroutines.FlowPreview +import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.reactive.asFlow import kotlinx.coroutines.reactive.awaitFirstOrNull import kotlinx.coroutines.reactive.awaitSingle -import kotlinx.coroutines.reactive.flow.asFlow import kotlin.reflect.KClass /** @@ -114,12 +114,9 @@ suspend fun ReactiveSelectOperation.TerminatingSelect.awaitExists() /** * Coroutines [Flow] variant of [ReactiveSelectOperation.TerminatingSelect.all]. * - * Backpressure is controlled by [batchSize] parameter that controls the size of in-flight elements - * and [org.reactivestreams.Subscription.request] size. - * * @author Sebastien Deleuze * @since 2.2 */ -@FlowPreview -fun ReactiveSelectOperation.TerminatingSelect.flow(batchSize: Int = 1): Flow = - all().asFlow(batchSize) +@ExperimentalCoroutinesApi +fun ReactiveSelectOperation.TerminatingSelect.flow(): Flow = + all().asFlow() diff --git a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensionsUnitTests.kt b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensionsUnitTests.kt index 757506ef4..b48a6b1bf 100644 --- a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensionsUnitTests.kt +++ b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensionsUnitTests.kt @@ -18,7 +18,7 @@ package org.springframework.data.cassandra.core import io.mockk.every import io.mockk.mockk import io.mockk.verify -import kotlinx.coroutines.FlowPreview +import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.flow.toList import kotlinx.coroutines.runBlocking import org.assertj.core.api.Assertions.assertThat @@ -220,7 +220,7 @@ class ReactiveSelectOperationExtensionsUnitTests { } @Test // DATACASS-648 - @FlowPreview + @ExperimentalCoroutinesApi fun terminatingFindAllAsFlow() { val spec = mockk>()