RxJava API轮询flatMap异常:FAILED状态未返回指定Observable
API轮询中FAILED状态未正确触发Error分支的问题排查与修复
问题背景
我有如下API轮询需求:
- API返回的
status可选值为"PENDING"、"FAILED"或"SUCCESSFUL" - 需持续轮询直到API返回"SUCCESSFUL"或"FAILED"
- 若返回"SUCCESSFUL",需调用另一个API获取附加数据
实现代码如下:
private fun pollApi(): Observable<Response> { return Observable.interval(3L, TimeUnit.SECONDS) .startWith(0L) .flatMap { api.getData() .subscribeOn(schedulerProvider.io()) .toObservable() } } override fun pollForPaymentState(): Observable<MyState> { return pollApi() .subscribeOn(schedulerProvider.io()) .takeUntil { val status = it.status (status == "SUCCESSFUL") || (status == "FAILED") } .flatMap { response -> when (response.status) { "FAILED" -> { Observable.just(MyState.Error) } "SUCCESSFUL" -> { // fetch additional data getAdditionalData() .subscribeOn(schedulerProvider.io()) .map { additionalDataResponse -> MyState.Success(additionalDataResponse.data) } .toObservable() } else -> Observable.never() } } } sealed class MyState { object Error: MyState() data class Success(additionalData: String): MyState() }
当前代码在status为"PENDING"或"SUCCESSFUL"时运行正常,但当status变为"FAILED"时,本该返回Observable.just(MyState.Error),却执行了else分支返回Observable.never(),导致无值返回。
问题根源分析
出现该问题的核心原因是硬编码字符串匹配的脆弱性:
- 后端返回的
status字段可能存在大小写差异(比如实际返回"failed"而非"FAILED")、前后空格,或者拼写错误(比如"FAIL"),导致无法匹配到"FAILED"分支 - 原代码依赖字符串字面量做分支判断,没有处理异常情况,一旦状态值不符合预期就会进入else分支
另外,原pollApi方法使用flatMap可能引发并发请求问题:interval定时发射事件时,如果前一次API请求还未完成,flatMap会同时发起新请求,造成不必要的资源消耗。
修复方案
1. 用枚举类替代硬编码字符串(推荐)
定义枚举类统一管理状态值,从根源避免字符串匹配错误:
enum class ApiStatus { PENDING, FAILED, SUCCESSFUL }
2. 优化轮询逻辑与分支判断
修改pollForPaymentState方法,先将字符串状态转换为枚举,同时处理异常情况:
override fun pollForPaymentState(): Observable<MyState> { return pollApi() .subscribeOn(schedulerProvider.io()) // 将字符串状态转换为枚举,处理可能的不匹配场景 .map { response -> val status = try { ApiStatus.valueOf(response.status.trim().uppercase()) } catch (e: IllegalArgumentException) { // 未知状态默认视为FAILED,也可根据业务需求抛出异常 ApiStatus.FAILED } Pair(response, status) } // 仅在状态为SUCCESSFUL或FAILED时终止轮询 .takeUntil { (_, status) -> status == ApiStatus.SUCCESSFUL || status == ApiStatus.FAILED } .flatMap { (response, status) -> when (status) { ApiStatus.FAILED -> Observable.just(MyState.Error) ApiStatus.SUCCESSFUL -> { getAdditionalData() .subscribeOn(schedulerProvider.io()) .map { MyState.Success(it.data) } .toObservable() } // 此处不会触发,因为takeUntil已过滤PENDING状态 ApiStatus.PENDING -> Observable.never() } } }
3. 修复轮询并发问题
将pollApi中的flatMap替换为concatMap,确保前一次API请求完成后再发起下一次轮询,避免并发:
private fun pollApi(): Observable<Response> { return Observable.interval(3L, TimeUnit.SECONDS) .startWith(0L) .concatMap { // 用concatMap保证请求串行执行 api.getData() .subscribeOn(schedulerProvider.io()) .toObservable() } }
额外优化
为MyState.Success的构造函数参数添加val,方便外部访问数据:
data class Success(val additionalData: String): MyState()
内容的提问来源于stack exchange,提问作者Mehdi Satei
相关产品推荐
相关产品推荐

