java-项目反应堆通量的并行处理
作者:互联网
我刚开始从事反应堆或反应式编程的工作,所以我可能做错了什么.我正在努力建立一个执行以下任务的流程:
给定一个类Entity:
Entity {
private Map<String, String> items;
public Map<String, String> getItems() {
return items;
}
}
>从数据库读取实体(ListenableFuture< Entity> readEntity())
>对每个项目(布尔processItem(Map.Entry< String,String>项目))执行一些并行异步处理
>当所有完成的调用doneProcessing(void doneProcessing(boolean b))时
目前,我的代码是:
handler = this;
Mono
.fromFuture(readEntity())
.doOnError(t -> {
notifyError(“some err-msg” , t);
return;
})
.doOnSuccess(e -> log.info("Got the Entity: " + e))
.flatMap( e -> Flux.fromIterable(e.getItems().entrySet()))
.all(handler::processItem)
.consume(handler::doneProcessing);
它可以正常工作,但是handler :: processItem调用并非在所有项目上同时运行.我尝试同时使用io和异步SchedulerGroup以及各种参数使用dispatchOn和publishOn,但是调用仍然在一个线程上串行运行.
我究竟做错了什么?
另外,我敢肯定,总体而言,上述内容可以改善,因此任何建议都会受到赞赏.
谢谢
解决方法:
您需要另一个flatMap可以为每个单独的地图元素派生并加入计算:
Mono.fromFuture(readEntity())
.flatMap(v -> Flux.fromIterable(v.getItems().entrySet()))
.flatMap(v -> Flux.just(v)
.publishOn(SchedulerGroup.io())
.doOnNext(handler::processItem))
.consume(handler::doneProcessing);
标签:project-reactor,java 来源: https://codeday.me/bug/20191027/1942726.html