【发布时间】:2021-10-02 19:41:50
【问题描述】:
我正在尝试将 CompletableFuture<Optional<T>> 转换为 Flow<T?>。我正在尝试编写的扩展函数是
fun <T> CompletableFuture<Optional<T>>.asFlowOfNullable(): Flow<T?> =
this.toMono().map { (if (it.isPresent) it.get() else null) }.asFlow()
但它失败了,因为 asFlow() 不存在可空类型,AFAICT 基于其定义。
那么,如何将CompletableFuture<Optional<T>> 转换为Flow<T?>?
编辑 1:
这是我到目前为止的想法。感谢您的反馈。
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flowOf
import java.util.Optional
import java.util.concurrent.CompletableFuture
fun <T> Optional<T>.orNull(): T? = orElse(null)
fun <T> CompletableFuture<Optional<T>>.asFlowOfNullable(): Flow<T?> = flowOf(this.join().orNull())
仅供参考,就我而言,它使用 Axon 的 Kotlin 扩展 queryOptional,我现在可以这样写:
inline fun <reified R, reified Q> findById(q: Q, qgw: QueryGateway): Flow<R?> {
return qgw.queryOptional<R, Q>(q).asFlowOfNullable()
}
我将推迟一段时间使用上述模式创建评论作为允许反馈的答案。
编辑 2:
由于下面指出编辑 1 中的 asFlowOfNullable 会阻塞线程,所以我现在从 @Joffrey 开始:
fun <T> Optional<T>.orNull(): T? = orElse(null)
fun <T> CompletableFuture<Optional<T>>.asDeferredOfNullable(): Deferred<T?> = thenApply { it.orNull() }.asDeferred()
编辑 3:感谢 @Tenfour04 和 @Joffrey 提供的有用意见。 :)
【问题讨论】:
-
在 Kotlin 中,我们通常不使用
Flow来表示单个项目。使用返回值的简单挂起函数更自然。为什么你需要一个流程? -
因为我使用 Axon 的 Kotlin 扩展来调用
QueryGateway.query<R,Q>(query:Q): CompletableFuture<R>,其中只有一个项目或 null。 -
这解释了为什么你有一个 CompletableFuture,而不是为什么你想将它转换为 Flow 而不是挂起函数或 Deferred。
-
由于
join调用,您在编辑 1 中的代码将阻塞调用线程直到完成。
标签: kotlin reactive-programming kotlin-extension kotlin-coroutines