【发布时间】: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