【问题标题】:Kotlin Process Collection In Parallel?Kotlin 进程收集并行?
【发布时间】:2018-01-16 10:45:39
【问题描述】:

我有一组对象,我需要对其执行一些转换。目前我正在使用:

var myObjects: List<MyObject> = getMyObjects()

myObjects.forEach{ myObj ->
    someMethod(myObj)
}

它工作正常,但我希望通过并行运行 someMethod() 来加快它,而不是等待每个对象完成,然后再开始下一个对象。

在 Kotlin 中有没有办法做到这一点?也许是doAsyncTask 之类的?

我知道这在 asked over a year ago 的时候是不可能的,但是现在 Kotlin 有像 doAsyncTask 这样的协程我很好奇这些协程是否可以提供帮助

【问题讨论】:

  • 这个问题已经有一年半的历史了,在引入 Kotlin 协程之前就被问到了
  • launch() 协程可以完成这项工作
  • 另外,您可以随时使用 Java 并行流

标签: collections parallel-processing kotlin kotlinx.coroutines


【解决方案1】:

你可以使用RxJava来解决这个问题。

List<MyObjects> items = getList()

Observable.from(items).flatMap(object : Func1<MyObjects, Observable<String>>() {
    fun call(item: MyObjects): Observable<String> {
        return someMethod(item)
    }
}).subscribeOn(Schedulers.io()).observeOn(AndroidSchedulers.mainThread()).subscribe(object : Subscriber<String>() {
    fun onCompleted() {

    }

    fun onError(e: Throwable) {

    }

    fun onNext(s: String) {
        // do on output of each string
    }
})

通过订阅Schedulers.io(),一些方法被安排在后台线程上。

【讨论】:

    【解决方案2】:

    Java Stream 在 Kotlin 中使用简单:

    tasks.stream().parallel().forEach { computeNotSuspend(it) }
    

    但是,如果您使用的是 Android,如果您想要与低于 24 的 API 兼容的应用程序,则不能使用 Java 8。

    您也可以按照您的建议使用协程。但截至目前(2017 年 8 月),它还不是语言的一部分,您需要安装一个外部库。有很好的guide with examples

        runBlocking<Unit> {
            val deferreds = tasks.map { async(CommonPool) { compute(it) } }
            deferreds.forEach { it.await() }
        }
    

    请注意,协程是通过非阻塞多线程实现的,这意味着它们可以比传统的多线程更快。我在下面的代码中对 Stream 并行与协程进行了基准测试,在这种情况下,协程方法在我的机器上要快 7 倍。 但是您必须自己做一些工作以确保您的代码是“暂停”(非锁定),这可能非常棘手。在我的示例中,我只是调用 delay,这是一个 @ 987654325@库提供的函数。非阻塞多线程并不总是比传统的多线程快。如果您有许多线程除了等待 IO 之外什么都不做,它会更快,这就是我的基准测试正在做的事情。

    我的基准测试代码:

    import kotlinx.coroutines.experimental.CommonPool
    import kotlinx.coroutines.experimental.async
    import kotlinx.coroutines.experimental.delay
    import kotlinx.coroutines.experimental.launch
    import kotlinx.coroutines.experimental.runBlocking
    import java.util.*
    import kotlin.system.measureNanoTime
    import kotlin.system.measureTimeMillis
    
    class SomeTask() {
        val durationMS = random.nextInt(1000).toLong()
    
        companion object {
            val random = Random()
        }
    }
    
    suspend fun compute(task: SomeTask): Unit {
        delay(task.durationMS)
        //println("done ${task.durationMS}")
        return
    }
    
    fun computeNotSuspend(task: SomeTask): Unit {
        Thread.sleep(task.durationMS)
        //println("done ${task.durationMS}")
        return
    }
    
    fun main(args: Array<String>) {
        val n = 100
        val tasks = List(n) { SomeTask() }
    
        val timeCoroutine = measureNanoTime {
            runBlocking<Unit> {
                val deferreds = tasks.map { async(CommonPool) { compute(it) } }
                deferreds.forEach { it.await() }
            }
        }
    
        println("Coroutine ${timeCoroutine / 1_000_000} ms")
    
        val timePar = measureNanoTime {
            tasks.stream().parallel().forEach { computeNotSuspend(it) }
        }
        println("Stream parallel ${timePar / 1_000_000} ms")
    }
    

    我的 4 核计算机上的输出:

    Coroutine: 1037 ms
    Stream parallel: 7150 ms
    

    如果您取消注释掉两个 compute 函数中的 println,您将看到在非阻塞协程代码中,任务以正确的顺序处理,但不是使用 Streams。

    【讨论】:

      【解决方案3】:

      是的,这可以使用协程来完成。以下函数对集合的所有元素并行应用操作:

      fun <A>Collection<A>.forEachParallel(f: suspend (A) -> Unit): Unit = runBlocking {
          map { async(CommonPool) { f(it) } }.forEach { it.await() }
      }
      

      虽然定义本身有点神秘,但您可以按照您的预期轻松应用它:

      myObjects.forEachParallel { myObj ->
          someMethod(myObj)
      }
      

      并行映射可以类似的方式实现,见https://stackoverflow.com/a/45794062/1104870

      【讨论】:

      • 目前“CommonPool”无法访问 - 它在“kotlinx.coroutines”内部!
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2023-03-11
      • 2016-04-14
      • 1970-01-01
      • 1970-01-01
      • 2019-04-09
      • 2023-01-02
      • 2014-11-08
      相关资源
      最近更新 更多