committed by
Mark Paluch
parent
eda8a1fdb6
commit
da702c2a49
@@ -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 <T : Any> ReactiveSelectOperation.TerminatingSelect<T>.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 <T : Any> ReactiveSelectOperation.TerminatingSelect<T>.flow(batchSize: Int = 1): Flow<T> =
|
||||
all().asFlow(batchSize)
|
||||
@ExperimentalCoroutinesApi
|
||||
fun <T : Any> ReactiveSelectOperation.TerminatingSelect<T>.flow(): Flow<T> =
|
||||
all().asFlow()
|
||||
|
||||
@@ -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<ReactiveSelectOperation.TerminatingSelect<String>>()
|
||||
|
||||
Reference in New Issue
Block a user