ParallelFlowable的sequential()算子工作机制及执行逻辑确认
ParallelFlowable.sequential() 发射时机详解
你期望的前者完全正确——sequential()并不会等待所有并行任务完成才转换为Flowable,而是会在任意一个并行任务产出结果的第一时间,就把这个结果发射给下游的操作符。
具体来说,ParallelFlowable的sequential()方法本质是做实时合并:它会监听所有并行分支的输出,一旦某个分支有数据产生,就立刻将其推送到单线程的下游流中,完全不需要等其他并行任务结束。
结合你的场景来看:当你用ParallelFlowable执行5个长时间任务,接着调用sequential()再加上take(1)时,只要其中某一个任务最先完成并输出结果,take(1)就会立刻捕获到这个结果,同时还会触发整个流的终止(RxJava会自动取消剩余未完成的并行任务,避免不必要的资源消耗)。你完全不需要等剩下的4个任务执行完毕,就能拿到第一个完成的结果。
可以用一段简单的代码验证这个行为:
Flowable.range(1, 5) .parallel() .runOn(Schedulers.io()) .map(taskNum -> { // 模拟不同耗时的任务,任务1最快完成 Thread.sleep(taskNum * 1000); return "任务" + taskNum + "已完成"; }) .sequential() .take(1) .subscribe( result -> System.out.println("收到结果:" + result), error -> error.printStackTrace(), () -> System.out.println("流已终止") );
运行这段代码后,你会看到1秒后就打印出"收到结果:任务1已完成",随后流终止,剩下的任务2-5会被自动取消,不会继续执行。
内容的提问来源于stack exchange,提问作者Daksh
相关产品推荐
相关产品推荐

