RxJava批量调用完成后执行后续操作遇subscribe重复触发问题
问题原因分析
你遇到的问题核心在于对Flowable.concat的回调时机理解偏差:concat确实会按顺序依次执行每个Flowable,但每个Flowable的onNext事件都会触发你订阅时的回调。如果你的answerListCreated里有2个元素,那么concat后的流就会发射2次数据,自然会触发2次addCorrectsAnswersToQuestion调用。
你的需求是「所有调用完成后只执行一次后续操作」,这时候应该监听的是整个流的完成事件,而不是单个元素的发射事件。
解决方案
修改订阅逻辑,把后续操作放在onComplete回调中(当所有串联的Flowable都执行完毕时,才会触发这个回调):
val flowableList = answerListCreated.map { questionService.addAnswerToQuestion( questionId, it.id, MyUtils.getAccessTokenFromLocalStorage(context = mContext!!) ) } disposable = Flowable.concat(flowableList) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe( // 可选:如果需要处理单个addAnswerToQuestion的返回结果,在这里实现 { result -> /* 处理单个请求的返回值 */ }, // 错误处理:如果任何一个请求失败,会触发这里 { throwable -> // 比如Toast提示错误,或者日志记录 Log.e("RxJava", "添加答案失败", throwable) }, // 所有请求都成功完成后,执行这部分逻辑 { addCorrectsAnswersToQuestion(questionId) } )
如果你不需要处理单个请求的返回结果,可以简化为:
disposable = Flowable.concat(flowableList) .ignoreElements() // 忽略所有单个元素的发射,只关注流的完成/错误 .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe( { addCorrectsAnswersToQuestion(questionId) }, { throwable -> Log.e("RxJava", "批量添加答案失败", throwable) } )
补充说明
- 如果你不需要严格串行执行这些请求(比如并行发起调用也可以),可以把
concat换成merge,效率会更高,但两者的核心问题解决逻辑是一致的——都是监听整个流的完成事件。 - 不管用
Flowable还是Observable,这个回调时机的逻辑是通用的,切换类型后代码改动很小。
内容的提问来源于stack exchange,提问作者StuartDTO
相关产品推荐
相关产品推荐

