【问题标题】:Is Reactor's FlatMap Asynchronous?Reactor 的 FlatMap 是异步的吗?
【发布时间】:2023-02-13 16:22:01
【问题描述】:

我是反应式编程的新手,我正在通过 micronaut 框架和 kotlin 使用反应堆。我试图了解反应式编程的优点以及我们如何使用地图平面地图通过单核细胞增多症助焊剂.

我了解反应式编程的非阻塞方面,但如果对数据流的操作实际上是异步的,我会感到困惑。

我一直在阅读 FlatMap 并了解它们异步生成内部流,然后将这些流合并到另一个 Flux 而不维护顺序。我看过的许多图表都使它更容易理解,但当涉及到实际用例时,我有一些基本问题。

例子:

fun updateDetials() {
        itemDetailsCrudRepository.getItems()
            .flatMap { 
                customerRepository.save(someTransferObject.toEntity(it))
            }
    }

在上面的示例中,假设 itemDetailsCrudRepository.getItems() 返回特定实体的 Flux。 flatMap 操作必须将通量中的每个项目保存到另一个表中。 customerRepository.save() 将从 flux 中保存项目,我们通过数据类 someTransferObject 的实例获取所需的实体。

现在,假设 getItems() 查询返回了 10 个项目,我们需要在新表中保存 10 行。 flatMap 操作(将这些项目保存到新表中的操作)是一次(同步)应用于通量的每一项,还是所有的保存都是异步发生的?

我读到的一件事是如果subscribeOn(Scheduler.parallel())不是应用然后将 flatMap 操作一次应用于通量中的每个项目(同步地).这个信息对吗?

如果我的基础知识本身不正确,请纠正我。

【问题讨论】:

    标签: asynchronous reactive-programming project-reactor micronaut flatmap


    【解决方案1】:

    是的,您对反应式编程有很好的理解,并且您的问题是有效的。在您显示的示例中,flatMap 操作一次应用于通量中的每个项目,除非您使用 subscribeOn 指定不同的调度程序。默认情况下,当您订阅 FluxMono 时,反应管道中的操作将在调用 subscribe 的同一线程上执行。如果要异步应用 flatMap 操作,可以使用 subscribeOn 指定不同的调度程序。例如,您可以使用Scheduler.parallel()在不同的线程池上执行操作,这将允许并行执行多个操作。下面是一个示例,说明如何修改代码以使 flatMap 操作异步:

    fun updateDetails() {
        itemDetailsCrudRepository.getItems()
            .flatMap { 
                customerRepository.save(someTransferObject.toEntity(it))
            }
            .subscribeOn(Scheduler.parallel())
            .subscribe()
    }
    

    请注意,在此示例中,我们在反应管道的末尾添加了一个 subscribe 调用以启动操作的执行。还值得注意的是,即使您使用subscribeOn(Scheduler.parallel())Flux 中的项目顺序也可能不会保留。项目的顺序将由底层调度程序确定,并且不能保证项目将按照它们发出的顺序进行处理。但是,您可以使用其他运算符(例如 concatMapconcatMapSequential)来保留 Flux 中项目的顺序。

    【讨论】:

      猜你喜欢
      • 2018-07-07
      • 2021-07-14
      • 2020-10-28
      • 1970-01-01
      • 1970-01-01
      • 2018-07-04
      • 1970-01-01
      • 2023-01-12
      • 2018-11-07
      相关资源
      最近更新 更多