编程语言
首页 > 编程语言> > java-项目反应堆通量的并行处理

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