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。