From 35f62933b7d201db959de204aee5533a74af7a5e Mon Sep 17 00:00:00 2001 From: Sebastien Deleuze Date: Sat, 6 Apr 2019 16:27:44 +0200 Subject: [PATCH] DATACASS-648 - Add Flow extension to ReactiveSelectOperation. Original pull request: #159. --- .../core/ReactiveSelectOperationExtensions.kt | 16 ++++++++++++++++ ...ctiveSelectOperationExtensionsUnitTests.kt | 19 +++++++++++++++++++ 2 files changed, 35 insertions(+) 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 c49ac3ec9..a1bbf7e8b 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,8 +15,11 @@ */ package org.springframework.data.cassandra.core +import kotlinx.coroutines.FlowPreview +import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.reactive.awaitFirstOrNull import kotlinx.coroutines.reactive.awaitSingle +import kotlinx.coroutines.reactive.flow.asFlow import kotlin.reflect.KClass /** @@ -107,3 +110,16 @@ suspend fun ReactiveSelectOperation.TerminatingSelect.awaitCount(): */ suspend fun ReactiveSelectOperation.TerminatingSelect.awaitExists(): Boolean = exists().awaitSingle() + +/** + * 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.allAsFlow(batchSize: Int = 1): Flow = + all().asFlow(batchSize) 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 204ba91dc..605031cd8 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,11 +18,14 @@ 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.flow.toList import kotlinx.coroutines.runBlocking import org.assertj.core.api.Assertions import org.junit.Test import org.springframework.data.cassandra.domain.Person import org.springframework.data.cassandra.domain.User +import reactor.core.publisher.Flux import reactor.core.publisher.Mono /** @@ -213,4 +216,20 @@ class ReactiveSelectOperationExtensionsUnitTests { find.exists() } } + + @Test + @FlowPreview + fun terminatingFindAllAsFlow() { + + val spec = mockk>() + every { spec.all() } returns Flux.just("foo", "bar", "baz") + + runBlocking { + Assertions.assertThat(spec.allAsFlow().toList()).contains("foo", "bar", "baz") + } + + verify { + spec.all() + } + } }