【问题标题】:Reactor not honoring runOn after flatMap call在 flatMap 调用后反应堆不遵守 runOn
【发布时间】:2019-07-05 17:59:16
【问题描述】:

我正在使用 SpringData MongoDB Reactive Streams 驱动程序和执行以下操作的代码:

reactiveMongoOperations.changeStream(changeStreamOptions, MyObject.class)
    .parallel()
    .runOn(Schedulers.newParallel("my-scheduler", 4))
    .map(ChangeStreamEvent::getBody)
    .flatMap(o -> {
        reactiveMongoOperations.findAndModify(query, update, options, MyObject.class)
    })
    .subscribe(this::process)

我希望一切都在my-scheduler 中执行。实际发生的是flatMap 操作确实在my-scheduler 中执行,而我的process() 方法中的代码却没有。

有人可以解释为什么会这样 - 这是一个错误还是我做错了什么?如何让Flux 中定义的所有操作在同一个调度程序上执行?

【问题讨论】:

  • process 在哪个线程中执行?难道是mongo驱动改变了findAndModify中的线程?
  • @SimonBaslé,process 方法在通用线程(“Thread-nn”)中执行。如果我在flatMap() 调用之后添加另一个 runOn(),那么process 将在my-scheduler 中执行。不过,我认为我不需要这样做。

标签: java reactive-programming spring-data-mongodb project-reactor


【解决方案1】:

runOn() 指定用于运行并行线程的每个“轨道”的调度程序。它不会影响订阅者。

如果您想为订阅者指定一个调度程序,那么您应该指定在原始Flux 上使用subscribeOn()(在parallel() 调用之前)。

【讨论】:

    猜你喜欢
    • 2018-09-01
    • 1970-01-01
    • 2016-08-29
    • 2023-03-13
    • 1970-01-01
    • 2014-12-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多