如何在并行调用Mono<Void>处理方法时收集Flux<Message>的处理错误messageId
问题根因
你的原始写法报错是因为process(m)返回的是Mono<Void>类型,而onErrorXXX系列方法要求返回值和上游的泛型类型保持一致,String和Void类型不匹配因此无法通过编译。
修正后可运行方案
@Data class Message { String messageId; String content; } // 最终返回的Mono<List<String>>就是所有处理失败的消息ID集合 Mono<List<String>> failedMessageIds = Flux.just(new Message("A", "1"), new Message("B", "2")) // flatMap第二个参数可自定义并行度,默认值为Queues.SMALL_BUFFER_SIZE = 256 .flatMap(m -> process(m) // 处理成功时返回空的String类型Mono,不向下游发射任何元素 .then(Mono.<String>empty()) // 处理出错时捕获异常,返回当前消息的messageId .onErrorResume(e -> Mono.just(m.getMessageId())) ) // 收集所有出错的messageId到List .collectList();
逻辑说明
- 我们统一了
flatMap内部返回值的泛型为String,彻底解决了类型不匹配问题 - 处理成功的消息不会向下游传递任何内容,只有处理失败的消息才会向下游输出自己的
messageId flatMap本身自带并行执行能力,你可以通过调整第二个入参控制最大并行处理数量,完全满足并行调用的需求
内容的提问来源于stack exchange,提问作者Daniel S. Hatten
相关产品推荐
相关产品推荐

