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

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(),导致无值返回。


问题根源分析

出现该问题的核心原因是硬编码字符串匹配的脆弱性:

  1. 后端返回的status字段可能存在大小写差异(比如实际返回"failed"而非"FAILED")、前后空格,或者拼写错误(比如"FAIL"),导致无法匹配到"FAILED"分支
  2. 原代码依赖字符串字面量做分支判断,没有处理异常情况,一旦状态值不符合预期就会进入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 20:25:39