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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 16:05:12