RxJava Completable Concat未按预期串行执行,求错误原因解析
嘿,我来帮你搞明白为啥会出现这种情况~
问题根源:你的Completable是「热启动」的
你期望Completable.concat能串行执行两个任务,但实际并行,核心问题出在completableTwoSeconds()方法的实现上:
当你调用completableTwoSeconds()的时候,方法内部的CompletableFuture.supplyAsync(...)会立刻启动异步任务——不管这个返回的Completable有没有被订阅。
而Completable.concat的工作逻辑是:依次订阅每个Completable,等前一个执行完成后,再订阅下一个。但你的代码里,两个completableTwoSeconds()在concat被调用时就已经被执行了,两个异步任务直接同时跑起来,自然就变成并行执行了。
解决方案:把任务改成「冷启动」(订阅才执行)
要让任务在被订阅时才启动,你有两种常见的修改方式:
方式1:用Completable.defer()包装
defer会延迟Completable的创建,直到它被订阅的时候才会调用你的方法生成新的Completable,这样就能保证串行执行:
@Test public void completableConcatTest() { Completable.concat( Completable.defer(this::completableTwoSeconds), Completable.defer(this::completableTwoSeconds) ) .doOnCompleted(() -> System.out.println("Both completed.")) .await(); } private Completable completableTwoSeconds() { PublishSubject<Void> subject = PublishSubject.create(); CompletableFuture.supplyAsync(() -> { try { Thread.sleep(2000); } catch (InterruptedException e) { e.printStackTrace(); } subject.onCompleted(); return null; }); return subject.doOnCompleted(() -> System.out.println("I'm completed." + Instant.now())).toCompletable(); }
方式2:用RxJava内置的fromCallable(更简洁)
RxJava提供了Completable.fromCallable()方法,它天然就是冷启动的,不需要手动创建Subject和CompletableFuture,代码会更简洁:
@Test public void completableConcatTest() { Completable.concat( completableTwoSeconds(), completableTwoSeconds() ) .doOnCompleted(() -> System.out.println("Both completed.")) .await(); } private Completable completableTwoSeconds() { return Completable.fromCallable(() -> { try { Thread.sleep(2000); } catch (InterruptedException e) { e.printStackTrace(); } return null; }).doOnCompleted(() -> System.out.println("I'm completed." + Instant.now())); }
总结一下
RxJava的concat操作符依赖的是「冷」的Observable/Completable(也就是只有被订阅时才会执行任务的类型),而你原来的实现是「热」的(创建时就启动任务),这才导致了并行执行的问题。用defer或者RxJava内置的冷创建方法就能解决这个问题啦~
内容的提问来源于stack exchange,提问作者alxg2112
相关产品推荐
相关产品推荐

