【问题标题】:async{} inside flowasync{} 内部流
【发布时间】:2021-08-25 11:36:32
【问题描述】:

我可以在 kotlin flow 中使用 async{} 吗?

场景:在 API 调用之后,我得到了一个包含 200 个对象的列表,我需要解析(转换为 UIObject)。 我正在尝试并行处理此列表。 下面是伪代码:

 fun getUIObjectListFlow():Flow<List<UIObject>> {
    flow<List<UIObject>> {
        while (stream.hasNext()) {
            val data = stream.getData() // reading from an input stream. Data comes in chunk

            val firstHalfDeffered = async(Dispatchers.IO) { /* process first half of the list that data obj contains*/ }
            val secondHalfDeffered = async(Dispatchers.IO) { /*process second half of the list that data obj contains */ }
            val processedList = firstHalfDeffered.await() + secondHalfDeffered.await() // just pseudo code

            emit(processedList)
        }
    }
}

由于 async{} 需要协程范围(例如: someScope.async{} ),我怎样才能在 flow 中获得范围?有没有其他方法可以做到这一点?

此函数在存储库中,我从 viewmodel 中调用它。

谢谢

【问题讨论】:

  • 如果只返回单个项目(即列表),为什么要首先使用流?如果您计划异步发出 UIObject 项目,那么您上面的代码不会这样做 - 它会等待所有项目,然后一次发出所有项目。物品的顺序重要吗?
  • @broot :知道了。在实际代码中,我正在读取流并发出对象而不是 List。我可以看到我的示例代码与我的问题不符。我会更新伪代码。
  • 在这种情况下,乔佛里的答案是正确的。您只需将流程主体包含在 coroutineScope { ... } 中。

标签: android kotlin kotlin-coroutines kotlin-flow


【解决方案1】:

(最初问题的原始答案)

正如 @broot 在 cmets 中提到的,如果您想要生成单个项目(即使该单个项目是一个集合),则不需要 Flow&lt;T&gt;。 通常,您只需要一个 suspend 函数(或本例中的暂停代码)而不是返回 Flow 的函数。

现在,无论您是否保留单项流,您都可以使用coroutineScope { ... } 挂起函数来定义一个本地范围,您可以从中启动协程。这个函数做了一些事情:

  1. 它提供了启动子协程的范围
  2. 它会暂停,直到所有子协程都完成
  3. 它根据块中的最后一个表达式返回一个值(lambda 的“返回”值)

这是它的样子:

val uiObjects = coroutineScope { //this: CoroutineScope
    val list = getDataFromServer()
            
    val firstHalf = async(Dispatchers.IO) { /*process first half of the list */ }
    val secondHalf = async(Dispatchers.IO) { /*process second half of the list */ }
            
    // the last expression from the block is what the uiObjects variable gets
    firstHalf.await() + secondHalf.await()
}

编辑:鉴于问题更新,这里是一些更新的代码。您仍然应该使用 coroutineScope 为您的短期协程创建本地范围:

fun getUIObjectListFlow(): Flow<List<UIObject>> = flow<List<UIObject>> {
    while (stream.hasNext()) {
        val data = stream.getData() // reading from an input stream. Data comes in chunk

        val processedList = coroutineScope {
            val firstHalfDeffered = async(Dispatchers.IO) { /* process first half of the list that data obj contains*/ }
            val secondHalfDeffered = async(Dispatchers.IO) { /*process second half of the list that data obj contains */ }
            firstHalfDeffered.await() + secondHalfDeffered.await() 
        }
        emit(processedList)
    }
}

【讨论】:

  • @Joffery:你能看一下我的更新伪代码吗?该函数从流中读取,然后处理数据块并从中创建一个 UIObject 列表,然后发出。
  • 正如我所说,无论你是否保持流量,你都可以做同样的事情。你可以用同样的方式使用coroutineScope
猜你喜欢
  • 1970-01-01
  • 2018-11-26
  • 1970-01-01
  • 2011-12-11
  • 1970-01-01
  • 2014-05-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多