使用flatMapDelayError时Flux遇异常后积压任务未正确执行问题咨询
问题原因分析
核心问题出在flatMapDelayError的特性和groupBy的错误处理逻辑的冲突上:
flatMapDelayError的作用是延迟错误到所有内部Publisher执行完成后再向下游合并抛出异常,而不是让错误不会流向下游- 当第一个
flatMapDelayError处理完所有元素后,会把收集到的0对应的IllegalStateException向下游发射给groupBy groupBy操作符只要收到上游的onError信号,无论之前已经产生的分组还有多少元素没处理,都会立刻终止所有分组流并转发错误,导致分组后的collectList还没收集完该分组的全部正常元素就提前报错,所以归约逻辑完全没有执行的机会
修复方案
要实现「先处理完所有正常数据,最后再抛出异常」的需求,我们需要先把错误转换成普通流元素,避免流提前触发onError终止,等所有正常数据处理完成后,再统一抛出收集到的错误:
修改后代码
public class Example1FlatMapDelayError1 { // 定义容器类区分正常数据和错误 private static class ResultHolder { private final Integer value; private final Throwable error; private ResultHolder(Integer value, Throwable error) { this.value = value; this.error = error; } public static ResultHolder success(Integer value) { return new ResultHolder(value, null); } public static ResultHolder error(Throwable error) { return new ResultHolder(null, error); } public boolean isSuccess() { return error == null; } public Integer getValue() { return value; } public Throwable getError() { return error; } } public static void main(String[] args) { ReactorDebugAgent.init(); ReactorDebugAgent.processExistingClasses(); Flux.just(1, 2, 3, 4, 5, 6, 7, 8, 9, 0) // 先把所有元素转成普通容器,不触发onError .map(integer -> { if (integer == 0) { return ResultHolder.error(new IllegalStateException()); } return ResultHolder.success(integer); }) // 分流:收集错误,处理正常数据 .publish(flux -> { // 收集所有错误 Flux<Throwable> errors = flux.filter(holder -> !holder.isSuccess()) .map(ResultHolder::getError); // 处理正常数据的归约逻辑 Flux<Integer> processResults = flux.filter(ResultHolder::isSuccess) .map(ResultHolder::getValue) .groupBy(integer -> integer % 2 == 0 ? "Even" : "Odd") .flatMap(group -> { System.out.println("Grouping Key-->" + group.key()); return group .collectList() .doOnNext(integers -> System.out.println("Key Value -->" + integers.stream().map(String::valueOf).collect(Collectors.joining(",")))) .flatMap(integers -> Flux.fromIterable(integers).reduce((i1, i2) -> i1 * i2)); }); // 先返回所有处理结果,最后合并错误抛出 return processResults.concatWith(errors.flatMap(Mono::error)); }) .subscribe( out -> System.out.println("Final -> " + out), throwable -> throwable.printStackTrace() ); } }
执行结果说明
修改后代码会先输出:
Grouping Key-->Odd Key Value -->1,3,5,7,9 Final -> 945 Grouping Key-->Even Key Value -->2,4,6,8 Final -> 384
等所有正常数据归约完成后,再抛出IllegalStateException异常,完全符合预期。
内容的提问来源于stack exchange,提问作者schengalath
相关产品推荐
相关产品推荐

