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>()