【问题标题】:Reactor spring mongodb repository combine multiple results togetherReactor spring mongodb 存储库将多个结果组合在一起
【发布时间】:2020-03-10 18:03:55
【问题描述】:

我是响应式编程的新手,目前正在开发基于 spring webflux 的应用程序。我被几个问题困住了。

public class FooServiceImpl {

@Autowired
private FooDao fooDao;

@Autowired
private AService aService;

@Autowired
private BService bService;

public long calculateSomething(long fooId) {
    Foo foo = fooDao.findById(fooId); // Blocking call one

    if (foo == null) {
        foo = new Foo();
    }

    Long bCount = bService.getCountBByFooId(fooId); // Blocking call two
    AEntity aEntity = aService.getAByFooId(fooId);  // Blocking call three

    // Do some calculation using foo, bCount and aEntity
    // ...
    // ...

    return someResult;
}
}

这是我们编写阻塞代码的方式,它使用三个外部 API 调用结果(我们将其视为 DB 调用)。我正在努力将其转换为反应式代码,如果所有三个都变成单声道并且如果我订阅所有三个,外部订阅者会被阻止吗?

public Mono<Long> calculateSomething(long fooId) {
    return Mono.create(sink -> {
        Mono<Foo> monoFoo = fooDao.findById(fooId); // Reactive call one
        monoFoo.subscribe(foo -> {
            if (foo == null) {
                foo = new Foo();
            }

            Mono<Long> monoCount = bService.getCountBByFooId(fooId);  // Reactive call two

            monoCount.subscribe(aLong -> {
                Mono<AEntity> monoA = aService.getAByFooId(fooId);  // Reactive call three
                monoA.subscribe(aEntity -> {
                    //...
                    //...
                    sink.success(someResult);
                });
            });
        });
    };
  }

我看到有一个叫做 zip 的函数,但是它只对两个结果有效,那么有没有办法在这里应用它呢?

如果我们在 create 方法中订阅了一些东西会发生什么,它会阻塞线程吗?

如果您能帮助我,将不胜感激。

【问题讨论】:

  • 你不应该订阅,订阅者是发起调用的调用客户端,你的服务是发布者,客户端是订阅者。如果您希望获取值并对其进行转换,则应使用 flatMap,如果您希望进行阻塞调用,则应将它们放在自己的调度程序中,如此处所述 projectreactor.io/docs/core/release/reference/…
  • 你所做的每一个subscribe 都是阻塞的,所以永远不要订阅(几乎,有一些边缘情况)
  • 请在 youtube 上查找教程或其他内容以了解如何编写基本的 webflux 应用程序。
  • 非常感谢@ThomasAndolf。我会调查的。

标签: reactive-programming spring-webflux project-reactor reactor reactor-netty


【解决方案1】:

如果你给了我你希望你用这些值做的计算,我会更容易展示反应堆的方法。但是让我们假设您想从数据库中读取一个值,然后将该值用于另一件事。使用平面图并制作独特的 Flux 减少代码行数和复杂性,无需像其他人所说的那样使用 subscribe()。示例:

return fooDao.findById(fooId)
.flatmap(foo -> bService.getCountBByFooId(foo))
.flatmap(bCount -> aService.getAByFooId(fooId).getCount()+bCount);

【讨论】:

  • 非常感谢您展示它,所以如果我们想合并并转换 map()flatmap() 是正确的方法吗?
  • 还要在另一个线程上运行回调事件,我们应该使用Schedulers 对吗?这意味着如果我们不使用Scheduler 调用,调用者线程将在subscribe 处被阻塞,直到调用onNext 回调方法,不是吗?如果我错了,请纠正我。
  • 阅读更多内容发现 Reactor 提供了一个名为 publishOn 的方法,用于在不同的线程上处理阻塞 API。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-07-07
  • 2021-07-23
  • 1970-01-01
  • 1970-01-01
  • 2014-11-10
  • 2022-07-07
相关资源
最近更新 更多