You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 07:15:58