如何在Kotlin中构建从Completable到Single的RxJava串行执行链?
RxJava串行执行Completable后接Single的问题解决
问题场景
在Kotlin中使用RxJava时,需要构建一条以Completable开头、后续串行执行Single操作的链,但当前代码中performCompletableTask()与performSingleStringTask()会被同时调用,无法保证串行执行。原代码如下:
class TestRepository { fun performCompletableTask(): Completable { Timber.v("Migration performCompletableTask") return Completable.fromAction { // Perform some operation here for (i in 1..100) { Timber.v("$i ") } } } fun performSingleStringTask(): Single<String> { Timber.v("Migration performSingleStringTask") return Single.fromCallable { val stringBuilder = StringBuilder() for (i in 1..100) { stringBuilder.append(i).append(" ") } stringBuilder.toString() } } fun performSingleIntTask(): Single<Int> { Timber.v("Migration performSingleIntTask") return Single.fromCallable { var sum = 0 for (i in 1..100) { sum += i } sum } } } fun testMigration() { Timber.v("Performing migration...") val completable = testRepository.performCompletableTask() .andThen(testRepository.performSingleStringTask()) .flatMap { Timber.v("flatMap Single<String> result: $it") testRepository.performSingleIntTask() } val subscribe = completable .subscribe({ Timber.v("subscribe Single<Int> result: $it") Timber.v("Migration completed.") }, { error -> Timber.e(error, "Error occurred during migration: ${error.message}") }) }
问题原因
原代码中andThen(testRepository.performSingleStringTask())直接调用了performSingleStringTask()方法,导致该方法内的初始化代码(比如Timber.v("Migration performSingleStringTask"))立即执行,Single实例也提前创建,而非等待performCompletableTask()执行完成后才触发。
解决方案
使用andThen的重载版本,接受Callable<Single<T>>(Kotlin中可直接用lambda替代),让RxJava在Completable执行完成后才调用Callable,进而触发后续Single方法的执行,保证串行顺序。
修改后的testMigration函数代码如下:
fun testMigration() { Timber.v("Performing migration...") val completable = testRepository.performCompletableTask() // 用Callable延迟执行performSingleStringTask() .andThen(Callable { testRepository.performSingleStringTask() }) .flatMap { Timber.v("flatMap Single<String> result: $it") // 同理,若需延迟performSingleIntTask(),也可包裹进Callable Callable { testRepository.performSingleIntTask() } } val subscribe = completable .subscribe({ Timber.v("subscribe Single<Int> result: $it") Timber.v("Migration completed.") }, { error -> Timber.e(error, "Error occurred during migration: ${error.message}") }) }
关键说明
- 通过
Callable包裹方法调用,将performSingleStringTask()的执行时机延迟到performCompletableTask()完成之后,确保操作串行执行。 - 若后续的
performSingleIntTask()也存在提前执行问题,同样可以用Callable包裹,保证整条链的串行顺序。
内容的提问来源于stack exchange,提问作者Shashank Pednekar
相关产品推荐
相关产品推荐

