Kotlin如何组合存在依赖关系的两个CallbackFlow完成数据聚合
问题背景
- 已定义两个数据类:携带ID字段的
MySection、MyArticle,无需关注两个类内部其余字段的具体实现。 - 获取Section列表的Flow实现如下:
private fun getSectionFlow(): Flow<List<MySection>> = callbackFlow { val listCallBackFlow = object : SectionCallback<List<MySection>>() { override fun onSuccess(sections: List<MySection>) { trySend(sections) } override fun onError(errorResponse: ErrorResponse) { close(Exception(errorResponse.reason)) } } provider.getSections( listCallBackFlow ) awaitClose() }
- 获取Article列表的Flow需要传入Section ID作为参数,才能查询对应分类下的文章列表,实现如下:
private fun getArticleFlow(id: Long): Flow<List<MyArticle>> = callbackFlow { val listCallBackFlow = object : ArticleCallback<List<MyArticle>>() { override fun onSuccess(sections: List<MyArticle>) { trySend(sections) } override fun onError(errorResponse: ErrorResponse) { close(Exception(errorResponse.reason)) } } provider.getArticle( id, // 传入section id查询对应文章列表 listCallBackFlow ) awaitClose() }
- 最终需要将两类数据整合为
List<MySectionWithArticle>结构,对应数据类定义:
data class MySectionWithArticle ( val id: Long, val title: String, var myArticlesList: List<MyArticle> = emptyList() )
实现要求
getArticleFlow(id)的入参ID需要取自getSectionFlow()返回结果中所有MySection的ID值:遍历getSectionFlow()返回的List<MySection>,为每个MySection实例调用getArticleFlow(id)获取对应分类下的文章列表,再将MySection基础信息和对应文章列表组合为MySectionWithArticle实例,最终返回MySectionWithArticle列表。
实现方案
直接通过Flow内置操作符组合即可实现,无需编写多层嵌套回调,分两种场景选择对应写法:
场景1:单次拉取数据(回调仅触发一次)
如果你的接口是一次性返回结果的类型,用flatMapLatest搭配协程并发请求即可,会保留Section原有顺序:
fun getSectionWithArticlesFlow(): Flow<List<MySectionWithArticle>> = getSectionFlow() .flatMapLatest { sectionList -> flow { // 开启协程域并发拉取所有分类下的文章 val assembledList = coroutineScope { sectionList.map { section -> async { // 取每个文章Flow的第一个返回结果即可 val articles = getArticleFlow(section.id).first() MySectionWithArticle( id = section.id, title = section.title, myArticlesList = articles ) } }.awaitAll() // 等待所有请求完成后返回组装结果 } emit(assembledList) } }
调用时如果只需要一次性拿结果,在协程内直接执行val result = getSectionWithArticlesFlow().first()即可。
场景2:响应式数据(数据源更新时自动推送新结果)
如果你的getArticleFlow是观察类数据源(比如Room数据库观察、回调会多次触发数据更新),改用combine操作符实现,任意分类下的文章更新时都会自动推送最新的组装列表:
fun getSectionWithArticlesFlow(): Flow<List<MySectionWithArticle>> = getSectionFlow().flatMapLatest { sectionList -> if (sectionList.isEmpty()) { flowOf(emptyList()) } else { // 为每个Section创建对应的文章观察Flow val articleFlows = sectionList.map { section -> getArticleFlow(section.id).map { articles -> section to articles } } // 组合所有Flow,每次数据更新时重新组装列表 combine(articleFlows) { resultPairs -> resultPairs.map { (section, articles) -> MySectionWithArticle( id = section.id, title = section.title, myArticlesList = articles ) } } } }
注意事项
你当前实现的callbackFlow中awaitClose()内部是空的,如果provider提供了取消回调注册的方法,一定要在awaitClose{}块内调用,否则会出现内存泄漏问题。
内容的提问来源于stack exchange,提问作者Pinkal Mistry
相关产品推荐
相关产品推荐

