<span class=“js_title_inner“>Coroutine 101: From Zero to Hero</span>

在日常生活中,当你在电话等待接通时,你可能会顺便查看电子邮件;在煮咖啡的同时,你也可能在做早餐;而开车时,你或许会听广播电台。
同样地,在编写软件时,我们也需要并行执行多项任务。例如同时发起两三个网络请求,与此同时不断更新界面,以显示每个请求的进度。
在 Kotlin 中,我们可以使用协程 (coroutines) 来实现多任务并行 (concurrency)。协程的用法多种多样,在本文中,我们只关注关于协程最实用的核心概念。
1. One Bot, One Thing at a Time
Steven 大部分时间都在为他的机器人 Bot 进行工程设计和管理,Bot 专门负责建筑项目。一天,Steven 和 Bot 接到一个新建筑施工的工作单。他们为这个项目准备的任务清单包括三项:
为地基砌砖
安装窗户
安装门
Steven 和 Bot 手头有一大堆砖,但窗户和门是针对每个项目单独订购的,必须从仓库调货。他们来到工地,准备开始工作。事情的发展如下:
在 t = 0 时等窗户:Bot 给仓库打电话订购窗户。送窗户花了 1.25 秒,在等待期间,它坐在人行道边闲得打转
在 t = 1.25 时等门:Bot 又给另一个仓库打电话订购门。送门花了 0.75 秒,它又坐在人行道边无聊地等门送来
在 t = 2 时砌砖:然后,它开始砌砖。这一部分花了 1 秒
在 t = 3 时安装窗户:接着,它花了 1 秒安装了窗户
在 t = 4 时安装门:最后,它花了 1 秒安装了门
在 t = 5 时完成任务
我们列出 Bot 做的工作时间线,总共花了 5 秒 (1.25 + 0.75 + 1 + 1 + 1)。

项目完成后,Steven 清点了所有任务,又摇了摇头:“这个项目花了太长时间。那些送货也太慢了。客户对我们这么久才交付也很不满意。我们能做些什么来加快进度呢?”
1.1 单线程阻塞 (Single-Threaded Blocking)
就像 Steven 的建筑项目一样,当我们的 Kotlin 代码一次只能做一件事,尤其是在等待诸如网络请求等慢操作时,也会效率低下。让我们用 Kotlin 代码来玩转吧。
首先,我们来创建一个枚举类 Product,表示可以从仓库订购的产品类型:
门:
DOORS("doors", 750),第一个参数是描述,第二个参数表示 750 毫秒的送货时间窗户:
WINDOWS("windows", 1_250),第一个参数是描述,第二个参数表示 1250 毫秒的送货时间
此外编写一个下单函数 order()和一个执行函数 perform()。我们使用 Thread.sleep() 来模拟卡车运送产品所需的时间。由于每种产品位于不同的仓库,它们的送货时间也各不相同。Thread.sleep() 接受一个 Long 类型的参数,表示要暂停的时长,单位是毫秒,比如要等待 1.25 秒钟,我们可以传入 1_250。
enum class Product(val description: String, val deliveryTime: Long) {
DOORS("doors", 750),
WINDOWS("windows", 1_250)
}
fun order(item: Product): Product {
println("ORDER EN ROUTE >>> The ${item.description} are on the way!")
Thread.sleep(item.deliveryTime)
println("ORDER DELIVERED >>> Your ${item.description} has arrived.")
return item
}
fun perform(taskName: String) {
println("STARTING TASK >>> $taskName")
Thread.sleep(1_000)
println("FINISHED TASK >>> $taskName")
}有了上面这些函数,我们就可以对 Steven 最近的项目进行建模,首先订窗户,然后订门,接着砌砖,再安装窗户,最后安装门。
fun main() {
val duration: Long = measureTimeMillis {
val window = order(Product.WINDOWS)
val door = order(Product.DOORS)
perform("laying bricks")
perform("installing ${window.description}")
perform("installing ${door.description}")
}
println("\nIt takes ${duration/1000.0} seconds to finish.")
}运行下面函数的产出如下,整个流程花了 5 秒 (精确值会比 5 多一点,考虑到程序运行并打印出来要花费的时间)。
ORDER EN ROUTE >>> The windows are on the way!
ORDER DELIVERED >>> Your windows has arrived.
ORDER EN ROUTE >>> The doors are on the way!
ORDER DELIVERED >>> Your doors has arrived.
STARTING TASK >>> laying bricks
FINISHED TASK >>> laying bricks
STARTING TASK >>> installing windows
FINISHED TASK >>> installing windows
STARTING TASK >>> installing doors
FINISHED TASK >>> installing doors
It takes 5.018 seconds to finish.就像 Steven 的建筑工作一样,这段代码完成任务的方式非常缓慢, Steven 不满意现状,看看他会提出哪些想法来提高工作效率。
1.2. 协程和并发 (Coroutines and Concurrency)
晚上 Steven 百无聊赖,突然想起之前和朋友玩三国杀里面的桥段,他选张飞连续出杀,朋友选赵云连续出闪。

如果将“张飞杀-赵云闪-张飞杀-赵云闪-张飞杀”出招过程写成 Kotlin 代码,应该长成下面的样子。
fun main() {
println("张飞 ➡️ ♠️A 杀")
println("赵云 ➡️ ♥️2 闪")
println("张飞 ➡️ ♣️3 杀")
println("赵云 ➡️ ♦️4 闪")
println("张飞 ➡️ ♠️5 杀")
}运行函数结果如下。
张飞 ➡️ ♠️A 杀
赵云 ➡️ ♥️2 闪
张飞 ➡️ ♣️3 杀
赵云 ➡️ ♦️4 闪
张飞 ➡️ ♠️5 杀上面代码“张飞杀-赵云闪-张飞杀-赵云闪-张飞杀”看上去很乱,最好是能把“张飞出杀”和“赵云出闪”的代码分开,如下所示。
fun main() {
println("张飞 ➡️ ♠️A 杀")
println("张飞 ➡️ ♣️3 杀")
println("张飞 ➡️ ♠️5 杀")
println("赵云 ➡️ ♥️2 闪")
println("赵云 ➡️ ♦️4 闪")
}但是运行函数结果改变了输出顺序 (张飞连出三个杀,赵云连出两个闪),不符合实际。
张飞 ➡️ ♠️A 杀
张飞 ➡️ ♣️3 杀
张飞 ➡️ ♠️5 杀
赵云 ➡️ ♥️2 闪
赵云 ➡️ ♦️4 闪现在问题来了,如何能将“张飞出杀”和“赵云出闪”的代码写在一起,而且能确保输出顺序正确?答案是用协程 (coroutine)!
使用协程编写的代码就像接力,一个协程可以做一些工作,然后“击掌换人”,让另一个协程运行一段时间。执行路径可以在各个协程之间交替,就像下图这样:

上面的代码演示了协程的精髓:执行路径可以在不同函数的各个部分之间来回切换。当代码以这种方式编写时,我们称这些任务是并发 (concurrency) 运行的。

准备好创建我们的第一个协程了吗?我们可以通过调用一种称为协程构建器 (coroutine builder) 的特殊函数来构建新的协程。我们使用的第一个协程构建器名为 runBlocking(),这个函数接受一个 lambda 函数,这是一个特殊的函数,称为挂起函数 (suspending function),大括号里面包含了该协程将要执行的代码。
import kotlinx.coroutines.runBlocking
fun main() {
runBlocking {
println("张飞 ➡️ ♠️A 杀")
println("张飞 ➡️ ♣️3 杀")
println("张飞 ➡️ ♠️5 杀")
println("赵云 ➡️ ♥️2 闪")
println("赵云 ➡️ ♦️4 闪")
}
}运行函数结果如下。
张飞 ➡️ ♠️A 杀
张飞 ➡️ ♣️3 杀
张飞 ➡️ ♠️5 杀
赵云 ➡️ ♥️2 闪
赵云 ➡️ ♦️4 闪奇怪,结果和没加runBlocking()的结果一样。这时候需要对“张飞出招”和“赵云出招”各加一个协程构建器 launch(),代码如下
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
fun main() {
runBlocking {
launch {
println("张飞 ➡️ ♠️A 杀")
println("张飞 ➡️ ♣️3 杀")
println("张飞 ➡️ ♠️5 杀")
}
launch {
println("赵云 ➡️ ♥️2 闪")
println("赵云 ➡️ ♦️4 闪")
}
}
}运行函数结果如下。
张飞 ➡️ ♠️A 杀
张飞 ➡️ ♣️3 杀
张飞 ➡️ ♠️5 杀
赵云 ➡️ ♥️2 闪
赵云 ➡️ ♦️4 闪真的要疯了,运行的结果还是一样。在 Kotlin 中,如果我们希望协程“击掌换人”,就必须放一个挂起点。一般来说,这会发生在它调用挂起函数的时候。为了演示这一点,我们来更新一下代码,使其在两人每次出招后都调用一个名为 yield() 的函数。
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
fun main() {
runBlocking {
launch {
println("张飞 ➡️ ♠️A 杀")
yield()
println("张飞 ➡️ ♣️3 杀")
yield()
println("张飞 ➡️ ♠️5 杀")
}
launch {
println("赵云 ➡️ ♥️2 闪")
yield()
println("赵云 ➡️ ♦️4 闪")
yield()
}
}
}运行函数结果如下,现在两人的出招顺序终于正确了。
张飞 ➡️ ♠️A 杀
赵云 ➡️ ♥️2 闪
张飞 ➡️ ♣️3 杀
赵云 ➡️ ♦️4 闪
张飞 ➡️ ♠️5 杀1.3 挂起函数 (Suspending Function)
让我们更新代码,使张飞或赵云每次出招时都先回回血。所以,与其直接调用 yield(),我们可以先创建一个新函数 restore(),它先打印 "续一秒!",然后再调用 yield()。如果你把它声明为一个普通函数,就会遇到编译错误 (见下划线标识处)。
fun restore() {
println("续一秒!")
yield()
}错误产生的原因在于,挂起函数只能由另一个挂起函数来调用!换句话说,普通函数只能调用普通函数,而挂起函数既可以调用普通函数,也可以调用其他挂起函数。

为了解决这个错误,我们只需在该函数前加上 suspend 修饰符,如下所示:
suspend fun restore() {
println("续一秒!")
yield()
}下面是在 main 函数中将原来的 yield() 调用替换为 restore() 的代码:
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
fun main() {
runBlocking {
launch {
println("张飞 ➡️ ♠️A 杀")
restore()
println("张飞 ➡️ ♣️3 杀")
restore()
println("张飞 ➡️ ♠️5 杀")
}
launch {
println("赵云 ➡️ ♥️2 闪")
restore()
println("赵云 ➡️ ♦️4 闪")
restore()
}
}
}运行函数结果如下,一切正常。
张飞 ➡️ ♠️A 杀
续一秒!
赵云 ➡️ ♥️2 闪
续一秒!
张飞 ➡️ ♣️3 杀
续一秒!
赵云 ➡️ ♦️4 闪
续一秒!
张飞 ➡️ ♠️5 杀在继续之前,回顾一下本节的主要概念:
协程可以相互并发运行,它们的执行可以被挂起,以便给其他协程运行的机会。
挂起函数能够挂起运行它们的协程。它们只能由其他挂起函数来调用。
runBlocking()会创建一个协程。位于其后的代码都要等到这个协程运行完毕后才会执行。launch()也会创建一个协程,位于其后的代码会在协程启动后立即执行。
虽然把协程比作接力很有趣,但 Steven 还有一些大型项目在等着他,他需要找到更高效的方法来完成它们!让我们看看接下来会发生什么。
2. One Bot, Two Things at a Time
从三国杀回到施工项目,Steven 思考着如何完成的更高效。正好他接到了另一个建筑项目的工单。Steven 心想“这一次每当我们下单订购物资,不要只坐等送货,让我们先去做下一个可做的任务。” 计划完 Steven 就让 Bot 马上动手。那天的流程如下:
在 t = 0 时:Bot 先给仓库打电话,订购窗户
在 t = 0 时:Bot 又给另一个仓库打电话,订购门
在 t = 0 时,Bot 随即开始砌砖 (用了 1 秒):
当还在砌砖时,第一辆送货车送来了门,花了 0.75 秒
砌完最后一块砖后过一会,第二辆送货车又送来了窗户,花了 1.25 秒
由于 1, 2, 3 同时开始,最后窗户送来后总共只花了 1.25 秒
在 t = 1.25 时:Bot 首先安装窗户,花了 1 秒
在 t = 2.25 时:Bot 最后安装门,花了 1 秒
在 t = 3.25 时完成任务
我们列出 Bot 做的工作时间线,总共花了 3.25 秒 (1.25 + 1 + 1)。

这个项目完成后,Steven 感到非常高兴。通过先做一部分工作、暂时放下,去处理其他任务,然后再回到最初的工作,他们节省了大量时间。Bot 坐在人行道边等待的时间少了许多!
2.1 模拟施工现场
现在已经了解了协程、协程构建器、挂起函数和挂起点的概念,接下来把这些知识带回到 Steven 和 Bot 身上。在小节 1.1 中,我们介绍了一个用于订购窗户和门等物资的函数 order(),一个用于完成任务的函数 perform()。现在,我们将把它改造为挂起函数,并用协程库中的挂起函数 delay() 来替换掉 Thread.sleep()。
import kotlinx.coroutines.delay
suspend fun order(item: Product): Product {
println("ORDER EN ROUTE >>> The ${item.description} are on the way!")
delay(item.deliveryTime)
println("ORDER DELIVERED >>> Your ${item.description} have arrived.")
return item
}再次详细解释 order() 函数:
作用:模拟向仓库下单并等待送货。
第一行打印
"ORDER EN ROUTE >>> …":表示订单已发出。delay(item.deliveryTime):根据item(窗户或门)在Product枚举中定义的送货时长挂起当前协程,期间让出线程给其他协程跑。第二行打印
"ORDER DELIVERED >>> …":表示货物送达。返回值:将收到的
item再传回调用方,以便后续取用(如安装时能拿到它的description)
suspend fun perform(taskName: String) {
println("STARTING TASK >>> $taskName")
delay(1_000)
println("FINISHED TASK >>> $taskName")
}再次详细解释 perform() 函数:
作用:模拟执行某个施工任务(如砌砖、安装)。
打印开始/结束:分别标记任务开始和结束。
delay(1_000):挂起 1 秒,模拟具体工作所需时间。
就像 Thread.sleep() 一样,delay() 函数也接受一个 Long 类型的参数,用来指定延迟的时长。两者的区别在于:
Thread.sleep()不会挂起协程。它只是简单地阻塞当前线程指定的时间段,因此无法在此期间让其他协程获得执行机会。该函数可以在普通函数或挂起函数中调用。delay()会挂起协程。这意味着协程可以在指定的时间段内“放下”当前工作,让其他协程在此期间运行。该函数只能在挂起函数内部调用。
在 Steven 最近的项目中,Bot 先给一个仓库打电话订购窗户,然后立刻又给另一个仓库打电话订购门。接着,他开始砌砖,却完全没有等这些物资送到。如果想在 Kotlin 中并发地处理多项任务,不能只把所有代码放到一个协程里,我们需要两个或更多的协程。因此,让我们改进代码,使 Bot 在等待送货的同时也能砌砖。
一种思路是:将每次 order() 调用都用 launch() 协程构建器包裹起来,然后再把所有的 perform() 调用也放到另一个 launch() 中。这样,两个 order() 调用就会在各自的协程里执行,与砌砖操作并发进行。不过,如果我们这么做,会遇到编译错误 (见下划线标识处)。
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val window = launch { order(Product.WINDOWS) }
val door = launch { order(Product.DOORS) }
launch {
perform("laying bricks")
perform("installing ${window.description}")
perform("installing ${door.description}")
}
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}造成此错误的原因在于,launch() 并不会返回 order() 的结果。它返回的是一个类型为 Job 的对象。这个 Job 对象很有用,但它无法让我们获取对 order() 调用的返回值。
相反,我们需要使用第三种协程构建器,名为 async()。这个构建器的用法与 launch() 非常类似,但它返回的不是 Job 对象,而是 Job 的一个子类型,称为 Deferred。该对象提供了一个名为 await() 的函数,可以让我们获取 order() 的执行结果。下面是如何使用它的示例。
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val window = async { order(Product.WINDOWS) }
val door = async { order(Product.DOORS) }
launch {
perform("laying bricks")
perform("installing ${window.await().description}")
perform("installing ${door.await().description}")
}
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}在这段代码中,在砌完砖 perform("laying bricks") 之后,我们会对两个 Deferred 对象,即延迟获取的窗户和延迟获取的门调用 await()。await() 是一个挂起函数,它会挂起当前协程,直到对应的 async() 协程执行完毕。

如果窗户在砌砖进行时就已经送达,那么窗户会在砌砖完成后立即安装。如果到砌砖结束时窗户还没有送到,那么由 launch() 创建的协程将会挂起,直到 order(Product.WINDOWS) 执行完毕为止。
整个流程花了大概 3.25 秒。订窗订门和砌砖在 0 秒处同时开始:订窗后挂起 1.25 秒,订门后挂起 0.75 秒,当砌砖在 1 秒处完成后:
先安装窗:窗在 1.25 秒处到达,但在砌砖完成后 (1 秒处) 还没到,因此到等到 1.25 处。等窗到位开始安装,因此在 1.25 - 2.25 秒区间完成。
后安装门:虽然门在 0.75 秒就已送达,但也要等到窗户安装后 (2.25 秒处) 才开始安装,因此在 2.25 -3.25 秒区间完成。

运行上面代码的结果如下。
ORDER EN ROUTE >>> The windows are on the way!
ORDER EN ROUTE >>> The doors are on the way!
STARTING TASK >>> laying bricks
ORDER DELIVERED >>> Your doors have arrived.
FINISHED TASK >>> laying bricks
ORDER DELIVERED >>> Your windows have arrived.
STARTING TASK >>> installing windows
FINISHED TASK >>> installing windows
STARTING TASK >>> installing doors
FINISHED TASK >>> installing doors
It takes 3.298 seconds to finish.
根据上面代码,我们得到了如下的协程结构:

2.2 不同施工现场
现场 1
如果将 window.await() 和 door.await() 顺序交换呢?
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val window = async { order(Product.WINDOWS) }
val door = async { order(Product.DOORS) }
launch {
perform("laying bricks")
perform("installing ${door.await().description}")
perform("installing ${window.await().description}")}
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}流程从 3.25 秒提高到 3 秒,为什么会变快?原因是订窗订门和砌砖在 0 秒处同时开始:订窗后挂起 1.25 秒,订门后挂起 0.75 秒,当砌砖在 1 秒处完成后:
先安装门:门在 0.75 秒就已送达,装门操作在砌砖完成后 (1 秒处) 就能立即进行,因此在 1 -2 秒区间完成。
后安装窗:窗在 1.25 秒到达,但只有在装门结束后 (2 秒处) 就绪,因此在 2 - 3 秒区间完成。
运行上面代码的结果如下。
ORDER EN ROUTE >>> The windows are on the way!
ORDER EN ROUTE >>> The doors are on the way!
STARTING TASK >>> laying bricks
ORDER DELIVERED >>> Your doors have arrived.
FINISHED TASK >>> laying bricks
STARTING TASK >>> installing doors
ORDER DELIVERED >>> Your windows have arrived.
FINISHED TASK >>> installing doors
STARTING TASK >>> installing windows
FINISHED TASK >>> installing windows
It takes 3.047 seconds to finish.
现场 2
如果将 async door 和 async window 的顺序交换呢?
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val door = async { order(Product.DOORS) }
val window = async { order(Product.WINDOWS) }
launch {
perform("laying bricks")
perform("installing ${door.await().description}")
perform("installing ${window.await().description}")
}
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}整个流程花了大概 3 秒。此改动只改变了打印门和窗在下单时的消息顺序,其他地方没有任何改变。
ORDER EN ROUTE >>> The doors are on the way!
ORDER EN ROUTE >>> The windows are on the way!
STARTING TASK >>> laying bricks
ORDER DELIVERED >>> Your doors have arrived.
FINISHED TASK >>> laying bricks
STARTING TASK >>> installing doors
ORDER DELIVERED >>> Your windows have arrived.
FINISHED TASK >>> installing doors
STARTING TASK >>> installing windows
FINISHED TASK >>> installing windows
It takes 3.038 seconds to finish.彻底理解并发里事件顺序的逻辑后,差不多该回到 Steven 的故事了,在此之前,让我们先回顾本节中的几个概念。
launch()构建器会返回一个Job对象,因此当你不需要从该协程中获取结果时,它是正确的选择。async()构建器会创建一个协程,并返回一个Deferred对象。你可以在该对象上调用await()函数来等待其结果。
Steven 和 Bot 通过并发执行工作取得了显著成效,但他们很快就会找到让项目更快完成的方法!让我们看看他们又在做些什么吧!
3. Two Bots, Two Things at a Time
又一个晚上 Steven 在家看 F1 比赛。当一辆赛车疾驰进维修站时,他惊叹于维修队的高效!几秒钟内,他们就将赛车顶起,精准更换车胎,并迅速为油箱加满燃油。之所以能如此迅速,是因为维修队有多名队员同时分工,各司其职。

Steven 心想:“如果我的施工团队里有多台机器人,项目会更快完成!”于是第二天,他开始制造更多机器人。在构建它们的过程中,他思考了项目中涉及的各种任务。有些任务需要搬运和砌砖这样的体力活,而另一些任务只需打电话订购材料并等待送货。因此,他将机器人分成了不同的小组:
第一组负责体力劳动和重物搬运,他称之为 Default 团队,它们承担主要的施工工作。
第二组负责输入 (Input) 输出 (Output) 的工作,称为 IO 团队,它们管场外设施之间的通信。
每个团队都有一名领班,根据当时的需求,负责将不同的任务分配给各个机器人。当下一个施工工单到来时,Steven 兴奋地想要试用他的新机器人团队!
在 t = 0 时:IO 团队给仓库打电话,订购窗户和门,并时刻关注它们的送达情况。
在 t = 0 时: Default 团队马上开始砌砖。
在 t = 0.75 时:门先送到,但 Default 团队还没砌完砖。由于砖先铺好才能安装门,这些门只好暂放一旁。
在 t = 1 时:砖已铺完,Default 团队一个机器人开始安装门。
在 t = 1.25 时:窗户送到。Default 团队另一个机器人随即准备好,窗户一到场就立刻安装。
在 t = 2.25 时完成任务
这种分工让项目比上一次节省了更多时间,总共花了 2.25 秒!(1.25 + 1)

Steven 兴奋极了!当然,并非所有事情都能同时进行。砌砖必须在安装窗户和门之前完成。然而,当他同时安装窗户和门时,项目却创纪录地提前完成了!如果我们想在 Kotlin 代码中实现同样的效果,就需要先了解
线程 (thread)、并发 (concurrency) 和并行 (parallelism)。
3.1 多线程并发 (Multi-Threaded Concurrency)
你在计算机上运行的任何程序,无论是你自己写的 Kotlin 程序、安装的应用,还是后台运行的服务,都在操作系统的线程上执行。如今大多数计算机都有多个处理器核心,因此可以同时处理多个线程。
就像 Steven 的机器人一样,我们可以用单线程在一段时间内做多件事;但如果计算机是多核的,我们也可以利用多个线程在同一时刻并行地完成多项任务。这就引出了一个重要的区分:
当代码的单一路径在两个或多个任务之间来回切换时,这些任务是在并发运行。
当存在多条执行路径,并且它们在同一瞬间各自执行不同任务时,这些任务是在并行运行。

到目前为止,我们的协程虽然能够并发执行代码,但始终都运行在同一个线程上。现在,让我们把这些协程部署到多个线程上,让它们真正并行运行吧!
就像 Steven 的每个团队都有一位领班负责将任务分配给团队中的机器人,在 Kotlin 中,我们也有不同的调度器 (Dispatcher),它们可以将协程分配到各自管理的线程上运行。如果想让某个调度器管理特定的协程,只需在协程构建器中传入该调度器即可。
3.2 调度器管理 (Dispatcher Management)
3.2.1 低效版
例如,我们可以使用 Dispatchers.IO 将产品订购任务分配给 IO 团队;而其他任务则可以使用 Dispatchers.Default 分配给 Default 团队。
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val window = async(Dispatchers.IO) { order(Product.WINDOWS) }
val door = async(Dispatchers.IO) { order(Product.DOORS) }
launch(Dispatchers.Default) {
perform("laying bricks")
perform("installing ${window.await().description}")
perform("installing ${door.await().description}")
}
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}此前,所有工作都在单个线程上执行,就是一台机器人完成了所有任务。但现在,通过将协程分配给 Dispatchers.IO 和 Dispatchers.Default,工作分别在三条不同的线程上进行:
由 IO 调度器管理的两条线程负责下单并监控物资送达。
由 Default 调度器管理的一条线程负责砌砖、安装窗户,安装门。
即便我们已经将工作分配给不同的团队,这些工作仍然花费了 3.25 秒才完成。为了进一步提升速度,我们需要并行安装窗户和门。
STARTING TASK >>> laying bricks
ORDER EN ROUTE >>> The windows are on the way!
ORDER EN ROUTE >>> The doors are on the way!
ORDER DELIVERED >>> Your doors have arrived.
FINISHED TASK >>> laying bricks
ORDER DELIVERED >>> Your windows have arrived.
STARTING TASK >>> installing windows
FINISHED TASK >>> installing windows
STARTING TASK >>> installing doors
FINISHED TASK >>> installing doors
It takes 3.291 seconds to finish.和以前一样,我们需要在希望与其他代码并发执行的代码周围,使用协程构建器。让我们用另一个 launch() 将最后两个 perform() 调用包裹起来。注意,由于必须先砌砖才能安装窗户和门,所以我们要等这部分工作完成后,才启动这些协程。
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val window = async(Dispatchers.IO) { order(Product.WINDOWS) }
val door = async(Dispatchers.IO) { order(Product.DOORS) }
launch(Dispatchers.Default) {
perform("laying bricks")
launch { perform("installing ${window.await().description}") }
launch { perform("installing ${door.await().description}") }
}
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}通过此改动,窗户和门得以同时安装,所有工作仅用大约 2.25 秒就完成了!
STARTING TASK >>> laying bricks
ORDER EN ROUTE >>> The doors are on the way!
ORDER EN ROUTE >>> The windows are on the way!
ORDER DELIVERED >>> Your doors have arrived.
FINISHED TASK >>> laying bricks
STARTING TASK >>> installing doors
ORDER DELIVERED >>> Your windows have arrived.
STARTING TASK >>> installing windows
FINISHED TASK >>> installing doors
FINISHED TASK >>> installing windows
It takes 2.289 seconds to finish.这段代码创建了 6 个不同的协程:其中 1 个由 runBlocking() 创建,2 个由 async() 创建,3 个由 launch() 创建。最终形成了如下的协程层级结构:

3.2.2 高效版
通过使用 withContext() 的函数,我们也可以仅用四个协程更高效地实现同样的效果,withContext() 的作用就是将工作交给另一个调度器。
到目前为止,我们已经把“订购产品和“安装产品”的代码分开了。

然而,“订购”与“安装”是紧密相关的操作,将它们在代码中放得更靠近会更有意义,同时又能保证每一步都由合适的团队负责。为此,我们无需再启动一个新的协程,只需使用 withContext()。
withContext() 函数允许我们在不启动新协程的情况下切换调度器。换句话说,它就像 IO 团队的某台机器人负责下单并等待送达;一旦货物到了,它再把窗户交给 Default 团队的机器人去完成实际的安装工作。代码示例如下:
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
launch(Dispatchers.IO) {
val windows = order(Product.WINDOWS)
withContext(Dispatchers.Default) { perform("install ${windows.description}")}
}
launch(Dispatchers.IO) {
val doors = order(Product.DOORS)
withContext(Dispatchers.Default) { perform("install ${doors.description}")}
}
launch(Dispatchers.Default) { perform("laying bricks") }
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}这段代码通过 launch() 创建了三个协程:一个负责处理窗户,一个负责处理门,另一个负责砌砖。在前两个协程中,产品的订购是在 Dispatchers.IO 管理的线程上进行的。但一旦产品到达就使用 withContext()切换调度器,让 perform() 在 Dispatchers.Default 管理的线程上执行。

通过此改动,我们创建了更少的协程,且相关工作 (例如订购窗户和安装窗户) 在代码中更加靠近。但这样也带来了一个问题!窗户和门本应在砌砖完成后才进行安装。如果查看输出,就会发现在砌砖工作尚未完成时就开始安装门了!
STARTING TASK >>> laying bricks
ORDER EN ROUTE >>> The doors are on the way!
ORDER EN ROUTE >>> The windows are on the way!
ORDER DELIVERED >>> Your doors have arrived.
STARTING TASK >>> install doors
FINISHED TASK >>> laying bricks
ORDER DELIVERED >>> Your windows have arrived.
STARTING TASK >>> install windows
FINISHED TASK >>> install doors
FINISHED TASK >>> install windows
It takes 2.289 seconds to finish.我们如何在开始安装门窗之前,等待砌砖任务完成呢?
还记得在调用 async() 构建器时,它会返回一个 Deferred 对象,我们可以在该对象上调用 await(),以挂起协程直到结果就绪。同样地,launch() 构建器会返回一个 Job 对象,其中包含一个名为 join() 的函数。与 await() 类似,join() 会挂起协程,直到 launch { … } 块中的代码执行完毕。让我们再次调整代码,这次确保在安装窗户和门之前,砌砖任务已经完成。
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val bricksJob = launch(Dispatchers.Default) { perform("laying bricks") }
launch(Dispatchers.IO) {
val windows = order(Product.WINDOWS)
bricksJob.join()
withContext(Dispatchers.Default) { perform("install ${windows.description}")}
}
launch(Dispatchers.IO) {
val doors = order(Product.DOORS)
bricksJob.join()
withContext(Dispatchers.Default) { perform("install ${doors.description}")}
}
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}上面代码的具体分析如下:
并发下单与砌砖
砌砖 (1 秒)、窗户下单 (1.25 秒)、门下单 (0.75 秒) 这三件事几乎同时开始,互不阻塞。
安装顺序由
join()控制两个安装任务 (窗户和门) 都会先执行
bricksJob.join(),保证它们在砌砖结束 (1 秒) 后再继续。对于门任务,因为它的下单只需 0.75 秒,所以到 1 秒时已完成,下单和等待砌砖结束都不再挂起;安装门的时间间隔为 1 秒 → 2 秒。
对于窗户任务,下单需 1.25 秒,切换到
join()时如果砌砖已结束 (1 秒),它会继续等到 1.25 秒才继续;安装窗户的时间间隔为 1.25 秒 → 2.25 秒。
线程切换
下单由
Dispatchers.IO线程池负责,避免阻塞 CPU 密集型线程;安装和砌砖由
Dispatchers.Default线程池处理,使得 CPU 密集型任务获得专用资源。
这段代码最终只创建了 4 个不同的协程:其中 1 个由 runBlocking() 创建,3 个由 launch() 创建。最终形成了如下的协程层级结构:

运行这段代码后,你会看到门和窗的安装要等到砌砖完成以后才会进行,整个过程耗时 2.25 秒。
STARTING TASK >>> laying bricks
ORDER EN ROUTE >>> The windows are on the way!
ORDER EN ROUTE >>> The doors are on the way!
ORDER DELIVERED >>> Your doors have arrived.
FINISHED TASK >>> laying bricks
STARTING TASK >>> install doors
ORDER DELIVERED >>> Your windows have arrived.
STARTING TASK >>> install windows
FINISHED TASK >>> install doors
FINISHED TASK >>> install windows
It takes 2.291 seconds to finish.Steven 和他的团队完美地完成了最近的项目,但在建筑行业中,你永远无法预料会有什么突发状况出现!让我们来看看他们即将面临哪些挑战吧。
4. Shit Happens
4.1 取消 (Cancellation)
4.1.1 全部取消
某天,Steven 和他的施工团队正忙着一个项目,这时客户打来电话。
“嘿,Steven,情况是这样的...” 客户开口说道,“我们要把整个业务搬到城里另一处地方。所以,你正在施工的那栋楼?我们不再需要了。就把整个项目取消吧。”
团队其实已经开工了,但既然客户不再需要,继续施工就没有意义了。有两台机器人正在等着窗户和门的送货,一台在砌砖。Steven 依次跑去告诉每台机器人这个消息。正在等待送货的那两台机器人立刻收到了通知,便收拾东西准备回家。
而正在砌砖的 Bot-3 因为戴着耳机全神贯注,一开始没注意到 Steven。直到最后一块砖砌完,他才抬头,看见 Steven 在示意大家准备回家,这才收拾好工具,结束了工作。

取消顶层协程 (cancel top-level coroutines)
就像在施工现场,有时需要取消某个协程任务。我们可以在传递给协程构建器的 lambda 内部调用名为 cancel() 的函数。让我们更新代码,在所有工作启动后取消该任务。
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val bricksJob = launch(Dispatchers.Default) { perform("laying bricks") }
launch(Dispatchers.IO) {
val windows = order(Product.WINDOWS)
bricksJob.join()
withContext(Dispatchers.Default) { perform("install ${windows.description}")}
}
launch(Dispatchers.IO) {
val doors = order(Product.DOORS)
bricksJob.join()
withContext(Dispatchers.Default) { perform("install ${doors.description}")}
}
cancel()
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}运行此代码后的输出如下:
STARTING TASK >>> laying bricks
ORDER EN ROUTE >>> The windows are on the way!
ORDER EN ROUTE >>> The doors are on the way!
FINISHED TASK >>> laying bricks
kotlinx.coroutines.JobCancellationException: BlockingCoroutine was cancelled
; job="coroutine#1":BlockingCoroutine{Cancelled}@7bedc48a通过调用 cancel(),我们得到了一个 JobCancellationException。就像在 Steven 的工地上一样,砌砖任务并未被中断。待会儿我们将了解这是为什么。现在,让我们更仔细地看看取消是如何工作的。
首先回顾该代码创建的协程层级结构:

属于同一层级的协程都存在于同一个协程作用域 (CoroutineScope) 中。实际上,CoroutineScope 是一个真正的接口,而 launch() 和 async() 这些协程构建器则是该接口的扩展函数,这让它们能够将新创建的协程绑定到相应的 CoroutineScope 上。
通过将协程组织到作用域中,Kotlin 可以追踪该作用域内的所有协程及其之间的父子关系。这样一来,如果某项工作被取消或出现异常,Kotlin 就能确保每个协程都得到妥善处理,而无需开发者手动管理这些情况。这个特性就称为结构化并发 (structured concurrency)。接下来,让我们看看当一个作业被取消时,结构化并发会如何发挥作用。
我们在 runBlocking() 的 lambda 内部调用了 cancel(),也就是在层级结构顶端的协程中。

多亏了结构化并发,当根协程被取消时,我们无需手动取消它的每一个子协程。相反,每个子协程都会自动收到取消信号;如果某个子协程恰好还有自己的子协程,它也会将该取消信号一并传递出去。

4.1.2 部分取消
某天,当施工队正忙着建造另一栋大楼时,客户打来电话,说:“我知道你们已经开始安装门了,但我们决定想要更通透的空间感。所以,不用再安装门了。我还是要这栋建筑,只是不要门。”
这一次,Steven 没有让所有机器人都停工,而是直接走到那个正在等门送达的机器人面前,向它发出了取消信号。窗户照常安装,建筑也顺利完工,只是没有门。
取消子协程 (cancel a child coroutine)
正如我们之前看到的,由于结构化并发,当你取消一个协程时,该协程以及它的所有子协程都会被取消。然而,这种取消不会影响它的父协程或兄弟协程。为了演示这一点,让我们像 Steven 的客户那样,取消门的那个任务。
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val bricksJob = launch(Dispatchers.Default) { perform("laying bricks") }
launch(Dispatchers.IO) {
val windows = order(Product.WINDOWS)
bricksJob.join()
withContext(Dispatchers.Default) { perform("install ${windows.description}")}
}
launch(Dispatchers.IO) {
val doors = order(Product.DOORS)
bricksJob.join()
cancel()
withContext(Dispatchers.Default) { perform("install ${doors.description}")}
}
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}运行此代码后的输出如下:
STARTING TASK >>> laying bricks
ORDER EN ROUTE >>> The windows are on the way!
ORDER EN ROUTE >>> The doors are on the way!
ORDER DELIVERED >>> Your doors have arrived.
FINISHED TASK >>> laying bricks
ORDER DELIVERED >>> Your windows have arrived.
STARTING TASK >>> install windows
FINISHED TASK >>> install windows
It takes 2.29 seconds to finish.这正如我们所期望的那样!可以看到门确实送到了,但它们并未被安装。其他一切都按计划进行,砖块被顺利铺好,窗户也已安装。因此,当我们取消一个协程时,该协程本身会被取消,如果它有子协程,它们也会被取消。但它的父协程和兄弟协程不会受到影响。

不过,取消并不是影响任务的唯一意外,Steven 和他的团队很快就要发现这一点了!
4.2 异常 (Exception)
那天,施工队又回到工地,开始另一个建筑项目。就在砌砖进行时,IO 团队的一台机器人打电话请求送门,却从仓库那头听到了令人意外的消息。
“抱歉,我们无法给你们送门了。你们的客户已经超出了预算,买不起更多的门。”电话那端传来这样的声音。
没有了门,也没有更多的资金,项目根本无法继续。那台机器人取消了当前工作,跑去告诉 Steven 发生了什么事。“看来我们只能放弃这个项目了,”Steven 说着,走向其他机器人,示意他们停止手头的工作。
有时候程序就是会遇到无法恢复的问题。如果异常没有被捕获,它最终会冒到调用栈顶端,导致整个应用崩溃。协程中也存在类似情形,不过它还会连带取消正在进行的工作。为演示这一点,我们改写之前的代码。
fun main() {
val duration: Long = measureTimeMillis {
runBlocking {
val bricksJob = launch(Dispatchers.Default) { perform("laying bricks") }
launch(Dispatchers.IO) {
val windows = order(Product.WINDOWS)
bricksJob.join()
withContext(Dispatchers.Default) { perform("install ${windows.description}")}
}
launch(Dispatchers.IO) {
val doors = order(Product.DOORS)
throw Exception("Out of money!")
bricksJob.join()
withContext(Dispatchers.Default) { perform("install ${doors.description}")}
}
}
}
println("\nIt takes ${duration / 1000.0} seconds to finish.")
}运行此代码后的输出如下:
STARTING TASK >>> laying bricks
ORDER EN ROUTE >>> The windows are on the way!
ORDER EN ROUTE >>> The doors are on the way!
ORDER DELIVERED >>> Your doors have arrived.
FINISHED TASK >>> laying bricks
java.lang.Exception: Out of money!尽管三个任务都已启动,但由于在下单门的任务内部抛出了未捕获的异常,它们都未能完成。默认情况下,协程中未捕获的异常会影响其作用域内的所有协程:
抛出异常的协程会取消它的所有子协程。
然后,它将异常向上传递给父协程,父协程接收到异常后也会取消它的所有子协程;这些子协程又会依此取消它们的子协程,如此递归。
这一过程会一直持续,直到异常传递到协程层级的顶端。

这种取消子协程并将异常在整个协程作用域内向上传播的行为是结构化并发的另一个特性。正如取消机制一样,它省去了我们为了确保所有协程正确关闭而需要手动完成的大量工作。
当夕阳映照在新建成的天际线上时,Steven 站在那里,欣赏着他和他的机器人施工团队取得的成果。最新的大楼不仅外观出色,而且在创纪录的时间内完工。他回想起几天前,Bot 还在路边低效地等待送货:他们已经走了多远啊!正是通过并发与并行地执行各项任务,他们用更少的时间建成了这些建筑,使客户满意度达到了历史新高。
当机器人陆续断电,Steven 也在夜幕中打起了盹,梦见有朝一日扩充他的机器人团队,去建造一座摩天大楼!

更多推荐

所有评论(0)