为了处理 HTTP 请求,我们必须将阻塞调用(例如 JDBC 调用)作为Mono/Flux基于进程的一部分。我们目前的计划是这样的:
// I renamed getSomething to processJaxrsHttpRequest
CompletionStage<String> processJaxrsHttpRequest(String input) {
return Mono.just(input)
.map(in -> process(in))
.flatMap(str -> Mono.fromCallable(() -> jdbcCall(str)).subscribeOn(fixedScheduler))
.flatMap(str -> asyncHttpCall(str))
.flatMap(str -> Mono.fromCallable(() -> jdbcCall(str)).subscribeOn(fixedScheduler))
.toFuture();
}
Run Code Online (Sandbox Code Playgroud)
wherefixedScheduler在 HTTP 请求中同时使用。
我们希望得到一些关于这种在相当数量的通量内处理块调用的策略的反馈。当然,我们知道,如果我们所有的请求都流经这些阻塞调用,那么我们最好不要使用反应器(除了公认的良好处理 API 之外)。
更新:感谢bsideup的回答。然而,我应该更具体地回答我的问题。
我的总体问题是,如果可以大量创建/订阅这些通量,如何有效地跨多个通量使用阻塞调用。我们尝试了建议的方法,但它会导致线程爆炸并很快出现 OOM。所以我们正在考虑使用共享调度程序。所以..这是我的问题。
fixedScheduler在我描述的情况下,您建议使用共享调度程序 ( ) 吗?如果没有,你能给我指点方向吗?Schedulers.newParallel("blocking-scheduler", maxNumThreads)?更新 2:刚刚深入研究Schedulers#newParallel并意识到这行不通,因为它“拒绝”阻塞呼叫。
非常感谢任何提示!
虽然subscribeOn这确实是处理阻塞调用的一种方法,并且您的用法没问题,但您也可以使用publishOn. 除非指定其他,否则
它将处理移至提供的:SchedulerpublishOn
CompletionStage<String> getSomething(String input) {
return Mono.just(input)
.map(in -> process(in)) // process must be non-blocking, or go after publishOn
.publishOn(Schedulers.boundedElastic())
.map(::jdbcCall)
.flatMap(str -> asyncHttpCall(str))
.publishOn(Schedulers.boundedElastic())
.map(::jdbcCall)
.toFuture();
}
Run Code Online (Sandbox Code Playgroud)
如您所见,您也可以继续使用异步调用。只需确保您没有阻止非阻塞调度程序(在该示例中,我publishOn再次使用 afterflatMap因为asyncHttpCall可能会从非阻塞调度程序完成)
| 归档时间: |
|
| 查看次数: |
2806 次 |
| 最近记录: |