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

Android Kotlin协程Flow中循环同步调用API报错问题排查

问题解决:Flow中forEach同步调用API避免AbortFlowException

问题背景

需求为在Android Kotlin协程Flow的forEach循环内同步调用API,等待所有调用完成且不阻塞主线程,但当前代码首次迭代可正常发射数据,后续迭代抛出错误:kotlinx.coroutines.flow.internal.AbortFlowException: Flow was aborted, no more elements needed。

错误原因

  1. flatMapLatest的特性冲突:flatMapLatest会在源Flow发射新值时,立即取消之前正在执行的下游Flow。如果getFirstList()频繁发射新数据,未完成的API调用会被强制中断,触发该异常。
  2. 冗余的Flow包装:repoApi.getCount()本身是挂起函数,无需额外包装成Flow,多余的Flow收集操作增加了复杂度和出错概率。
  3. 异常场景未处理:原代码未对data.items或datalist.Id的空值做防护,空指针可能间接导致Flow执行异常。

解决方案

核心调整点

  • 替换flatMapLatest为flatMapConcat(串行处理新值)或flatMapMerge(并行处理新值),避免旧的Flow任务被强制取消。
  • 直接调用挂起函数getCount(),移除冗余的Flow包装,简化同步调用逻辑。
  • 增加空值和异常处理,确保流程稳定性。

修改后的代码(串行调用)

fun getLatestData(): Flow<Model> = flow {
    emit(getFirstList())
}.flatMapConcat { data -> // 改用flatMapConcat,串行处理源Flow的新值,不取消旧任务
    flow {
        val responseList = mutableListOf<String>()
        data.items?.forEach { datalist ->
            datalist.Id?.let { id ->
                try {
                    // 直接调用挂起API,自动在IO线程执行(配合flowOn)
                    val countResponse = repoApi.getCount(id)
                    responseList.add(countResponse.check)
                } catch (e: Exception) {
                    e.printStackTrace()
                    // 可添加异常兜底逻辑,比如添加默认标记或跳过当前项
                }
            }
        }
        // 所有API调用完成后,统一发射结果
        emit(Model(
            datalist = data,
            responseListCount = responseList
        ))
    }.flowOn(Dispatchers.IO) // 将API调用切换到IO线程,不阻塞主线程
}

// 移除冗余的getCreatedDate函数,直接使用repoApi的挂起方法

@Headers(
    "Content-Type: application/json"
)
@GET("getCount/{Id}")
suspend fun getCount(@Path("Id") id: String): Response // 建议Id设为非空,避免空指针风险

可选:并行API调用优化

如果需要提升效率,可改用flatMapMerge结合async/await实现并行调用:

fun getLatestData(): Flow<Model> = flow {
    emit(getFirstList())
}.flatMapMerge { data ->
    flow {
        // 并行发起所有API请求
        val deferredList = data.items?.mapNotNull { datalist ->
            datalist.Id?.let { id ->
                async(Dispatchers.IO) {
                    try {
                        repoApi.getCount(id).check
                    } catch (e: Exception) {
                        e.printStackTrace()
                        null // 异常时返回null,后续过滤
                    }
                }
            }
        } ?: emptyList()
        
        // 等待所有请求完成,过滤异常项
        val responseList = deferredList.awaitAll().filterNotNull()
        
        emit(Model(
            datalist = data,
            responseListCount = responseList
        ))
    }
}

内容的提问来源于stack exchange,提问作者Android Dev

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 21:45:17