【发布时间】:2017-11-28 00:13:39
【问题描述】:
有没有办法获得与下面的代码 sn-p 相同的行为但使用协程?
更新代码 sn-p:
fun main(args: Array<String>) = runBlocking {
val executor = Executors.newFixedThreadPool(50)
log.info("Start")
val jobs = List(300) {
executor.submit {
log.info("worker #$it started")
sleep(1000L)
log.info("worker #$it done")
}
}
jobs.forEach { it.get() }
executor.shutdown()
log.info("All done!")
}
如何运行 300 个并行因子 == 50 的作业,但不创建 50 个实际线程?
更新 2:解决方案
再读一遍Coroutines Guide 后,我发现Fan-out example 正是我想要的。因此,我的示例如下所示:
fun produceTasks() = produce {
for (taskId in 1..300) {
send(
async(start = CoroutineStart.LAZY) {
delay(1000) // simulate long work
taskId
}
)
}
close()
}
fun launchWorker(index: Int, channel: ProducerJob<Deferred<Int>>) = launch {
channel.consumeEach {
val result = it.await()
log.info("Worker #$index done task #$result")
}
}
fun main(args: Array<String>) = runBlocking {
val tasks = produceTasks()
val workers = List(50) { launchWorker(it + 1, tasks) }
workers.forEach { it.join() }
log.info("Done")
}
【问题讨论】:
-
当您说“300 个并行因子 == 50 的作业”时,如果您认为这不是多个底层真实线程,那么“并行因子”是什么意思?
-
我的意思是作业应该以长度 50 排队,但使用轻量级协程(可能有大约 4 个真正的线程底层而不是 50 个真正的线程)。但是,如果我在下面的评论中写下类似的内容,那么所有 300 个工作/任务都是同时开始的。