前言

看过很多关于 Kotlin协程相关原理的文章,个人感觉比较模糊,可能是循循渐进的方式不太适合自己,还是觉得在一开始先阐述好某些已有结论,然后通过查看源码等方式来验证这种结论会比较好;

协程的挂起和恢复

Kotlin 的协程挂起和恢复,是在某个 挂起点 进行的,这个挂起点包括:

  • 执行其他挂起函数

  • 切换上下文(实际上也算是挂起函数)

在 IDEA 中在左侧可以看到一个箭头的图标时(如下图),就表明是一个挂起点;image-knld.png

既然是挂起,当前线程可以去干别的事情,而无需堵塞等待挂起函数的执行完毕,那肯定需要保存当前挂起函数的状态,就像栈帧一样保存当前函数执行位置以及局部变量等信息,方便在挂起点恢复时直接恢复挂起前的状态;

在编写 suspend 函数的时候,并看到没有相关的处理,就像是在写一个普通的函数,难道说 Kotlin 底层给这个函数栈帧额外的处理?

并没有,而是编译器为我们掩盖了这一切的处理,编译器将 suspend 函数编译成一个 Continuation 匿名对象,并在其内部实现了 状态机,以及通过成员字段来保存挂起点的信息以及挂起函数的执行结果:

  • 状态机:将挂起点之前的一系列操作视作为一个状态,当遇到挂起点时,将状态切换为执行完挂起点后的下一个状态;

  • 挂起:实际上就是当前状态下的方法结束,因此线程才能够去执行其他的协程;

  • 成员字段:当挂起时,也就是方法结束时,需要将当前函数的局部变量信息保存到成员字段中,挂起点恢复时,需要从成员字段中获取局部变量信息;

  • 恢复:挂起点的恢复实际上是再次调用该方法,因为状态机已经切换到下一个状态,所以执行的是挂起点后的代码;

协程的线程池

在 Kotlin 的协程中,通常会用到 Dispatcher 来切换线程,而这个 Dispatcher 实际上是一个线程池(除了Android 的 Dispatcher.Main 是一个 Handler)

上面提到的 Continuation 匿名对象在执行时,会被封装为一个实现 Runnable 接口的 Task,提交到线程池中执行,执行协程中对应状态的一系列逻辑;

在协程中的切换 Dispatcher 实际上是将这个协程在不同的线程池之间 流转

比如一个协程先在 Dispatcher.Main 执行一部分,切换到 Dispatcher.IO 执行另一部分,最终再切换到 Dispatcher.Main 执行剩余部分;结合上面说到的状态机,实际上就是分为三种状态:

  • 第一种状态在 Main 的线程池中执行(Android中是 Handler);

  • 第二种状态在 IO 的线程池中执行;

  • 第三种状态再次回到 Main 的线程池中执行;

协程的执行基于线程池,并能在不同的线程池之间流转,因此也说协程是一个轻量级的线程框架,框架为我们处理了线程池之间流转的代码,我们只需要通过 withContext 就能切换线程池;

提出问题

现在待探究的问题是,Kotlin 编译器如何将 suspend 函数编译为对应的 Continuation 匿名对象,以及协程是如何挂起和恢复的;

同时探究协程的一些特性是怎么实现的?

  • 同步方式写异步代码

  • 结构性并发(统一管理协程)

  • 异常的处理


编译器魔法

Kotlin 编译器对 suspend 修饰的函数进行编译处理,现在用非协程的异步写法和协程的写法进行对比,查看编译器到底处理了什么;

非协程的异步写法

如果不使用协程来编写异步代码,那肯定少不了 回调 这个东西;

假设需要新开一个线程来获取网络数据,对于网络数据的返回时间是不确定的,因此需要通过 回调 的方式来通知数据已经返回

 // 回调接口
 fun interface ResultCallback<T> {
     fun onResult(result: Result<T>)
 }
 ​
 // 模拟网络请求数据
 fun mockNetworkData(param: String): String {
     val tmp = param.repeat(2)
     Thread.sleep(1000)
     return "data: $tmp"
 }
 ​
 fun main() {
     val callback = ResultCallback<String> { res ->
         res.onSuccess {
             println("Get Data: $it")
         }
     }
     // 启动线程时需要给一个回调来通知数据返回
     val thread = thread {
         val data = mockNetworkData("123")
         callback.onResult(Result.success(data))
     }
     // 堵塞主线程,等待子线程执行完毕
     thread.join()
 }
 ​
 // 运行结果:
 // Get Data: data: 123123

如果对于某些需要等待数据返回才能请求其他数据的场景,就需要编写多层嵌套回调,整个代码层级会显得很复杂和臃肿;

协程的“异步”写法

协程的一大特点就是可以用同步的写法编写异步的代码,上面的代码可以替换为:

 // 模拟网络请求数据
 suspend fun mockNetworkData(param: String): String {
     val tmp = param.repeat(2)
     delay(1000)
     return "data $tmp"
 }
 ​
 fun main() {
     runBlocking {
         // 启动一个协程
         launch(Dispatchers.IO) {
             val data = mockNetworkData("123")
             println("Get Data: $data")
         }
     }
 }

当遇到等待数据返回才能请求其他数据的场景,协程只需要继续在下面调用函数即可,而不需要编写多层回调;

编译器的“魔法”

查看上面协程异步写法的 mockNetworkData 方法,反编译为 Java 代码,可以看到生成了很多本来没有的东西:

  • 本来没有参数的方法,新增加了一个 Continuation 类型的参数

  • 方法体内新增了标签 labelxxContinuationImpl 匿名类对象、switch 分支判断等内容;

 @Nullable
 public static final Object mockNetworkData(@NotNull String param, @NotNull Continuation $completion) {
    Object $continuation;
    label20: {
       if ($completion instanceof <undefinedtype>) {
          $continuation = (<undefinedtype>)$completion;
          if ((((<undefinedtype>)$continuation).label & Integer.MIN_VALUE) != 0) {
             ((<undefinedtype>)$continuation).label -= Integer.MIN_VALUE;
             break label20;
          }
       }
       $continuation = new ContinuationImpl($completion) {
          Object L$0;
          // $FF: synthetic field
          Object result;
          int label;
          @Nullable
          public final Object invokeSuspend(@NotNull Object $result) {
             this.result = $result;
             this.label |= Integer.MIN_VALUE;
             return TestKt.mockNetworkData((String)null, (Continuation)this);
          }
       };
    }
    Object $result = ((<undefinedtype>)$continuation).result;
    Object var5 = IntrinsicsKt.getCOROUTINE_SUSPENDED();
    String tmp;
    switch (((<undefinedtype>)$continuation).label) {
       case 0:
          ResultKt.throwOnFailure($result);
          tmp = StringsKt.repeat((CharSequence)param, 2);
          ((<undefinedtype>)$continuation).L$0 = tmp;
          ((<undefinedtype>)$continuation).label = 1;
          if (DelayKt.delay(1000L, (Continuation)$continuation) == var5) {
             return var5;
          }
          break;
       case 1:
          tmp = (String)((<undefinedtype>)$continuation).L$0;
          ResultKt.throwOnFailure($result);
          break;
       default:
          throw new IllegalStateException("call to 'resume' before 'invoke' with coroutine");
    }
    return "data " + tmp;
 }

Continuation 接口

新增的 Continuation 类型参数,其内部就有一个方法 resumeWith

  • result 是传入其他挂起函数的返回值,用于恢复协程的挂起点,实际上就是一个异步回调的方法;

  • result 可能是执行成功的返回值 ,也可能是执行出现的异常 Result.Failure

 public interface Continuation<in T> {
     public val context: CoroutineContext
     public fun resumeWith(result: Result<T>)
 }

编译器为我们自动传入了一个 异步回调,对于 非协程的异步写法,两种效果是一样的;

看下面这个例子,与启动一个线程的写法类似, suspen {...} 是一个协程体,传入一个异步回调 ContinuationstartCoroutine 是启动这个协程,当运行完毕后就会回调 Continuation#resumeWith 传回执行结果;

 suspend {
     delay(1000)
     "123" // 执行结果
 }.startCoroutine(object : Continuation<String> {
     override val context: CoroutineContext
         get() = EmptyCoroutineContext
     
     override fun resumeWith(result: Result<String>) {
         result.onSuccess {
             println("Get Data: $it")
         }
     }
 })
 // 运行结果
 // Get Data: 123

题外知识,Retrofit2 判断一个函数是否为 suspend 修饰,就是通过判断最后一个参数的类型是否为 Continuation 实现的;

ContinuationImpl 匿名类

接着反编译的内容,简单分析下 ContinuationImpl 匿名类:

 label20: {
     // Q1
     if ($completion instanceof <undefinedtype>) {
         $continuation = (<undefinedtype>)$completion;
         if ((((<undefinedtype>)$continuation).label & Integer.MIN_VALUE) != 0) {
             ((<undefinedtype>)$continuation).label -= Integer.MIN_VALUE;
             break label20;
         }
     }
     // Q2:创建匿名对象
     $continuation = new ContinuationImpl($completion) {
         Object L$0;
         // $FF: synthetic field
         Object result;
         int label;
         @Nullable
         public final Object invokeSuspend(@NotNull Object $result) {
             this.result = $result;
             this.label |= Integer.MIN_VALUE;
             return TestKt.mockNetworkData((String)null, (Continuation)this);
         }
     };
 }

Q1 判断了传入参数的类型是否为匿名类(虽说是匿名,JVM编译器还是会为匿名类创建对应的类),当已经创建过时会跳过再次创建;

Q2 创建了一个 ContinuationImpl 的匿名类对象:

  • 有多个字段,基本分为 resultlabelL$xx三种:

    • result:挂起点的运行结果;

    • L$xx:保存挂起点前的局部变量信息,不包括编译时常量,而是需要运行时才能确定下来的局部变量,比如 System.currentMillions();(该字段的数量与符合条件的局部变量的个数相同)

    • label:当前协程状态机的状态;

  • 其内部的 invokeSuspend 方法主要用于恢复挂起点,参数 $result 是挂起点函数的执行返回结果;其内部再次调用了 mockNetworkData 方法,传入的参数是当前的匿名类对象;

捋一下逻辑:

  1. 第一次调用挂起函数时,创建匿名类对象;

  2. 挂起点的挂起函数调用匿名类对象的 invokeSuspend 方法,将执行结果 $result 赋值给该匿名类对象,再次调用挂起函数;

  3. 再次调用挂起函数,将匿名类对象传入,走到 Q1 的判断分支成立,跳过创建匿名对象过程;

状态机和状态切换

里面的 switch 根据 continuation.label 进行判断实际上就是 状态机 以及 状态机的切换

 // 挂起点的执行结果,对于挂起点前是null,在挂起点后是非null
 Object $result = ((<undefinedtype>)$continuation).result;
 // 挂起状态的判断
 Object var5 = IntrinsicsKt.getCOROUTINE_SUSPENDED();
 String tmp;
 switch (((<undefinedtype>)$continuation).label) {
     case 0:
         // Q1
         ResultKt.throwOnFailure($result);
         tmp = StringsKt.repeat((CharSequence)param, 2);
         ((<undefinedtype>)$continuation).L$0 = tmp;
         ((<undefinedtype>)$continuation).label = 1;
         if (DelayKt.delay(1000L, (Continuation)$continuation) == var5) {
             return var5;
         }
         break;
     case 1:
         // Q2
         tmp = (String)((<undefinedtype>)$continuation).L$0;
         ResultKt.throwOnFailure($result);
         break;
     default:
         throw new IllegalStateException("call to 'resume' before 'invoke' with coroutine");
 }
 return "data " + tmp;

编译器根据挂起点 delay 函数划分为两个状态:

  • Q1:在调用挂起函数前保存局部变量数据 tmp ,并且将状态 label 切换至 1(下一个状态),最后调用 delay 挂起函数,并将匿名类对象传入,用于恢复;

  • Q2:挂起点恢复后,恢复局部变量数据 tmp,最后返回执行结果;

switch 分支前的赋值语句,有需要说明的点:

  • $result:这里是从匿名类对象中取出挂起点的执行结果,在挂起点前因刚创建的默认值为 null,在挂起点恢复后的值就是调用的挂起函数的返回值(Unit或者其他类型);

  • getCOROUTINE_SUSPENDED():这是用于判断一个挂起函数是否会 真正的挂起 的标志;(具体请看下一小节)

继续捋一下调用逻辑:

  1. 当首次调用 mockNetworkData 时,会创建匿名类对象 $continuation,此时状态 label 为 0,会进入 Q1 处执行;

    • 在挂起点前对局部变量 tmp 进行保存,并且将状态切换为 1;

    • 调用挂起函数 delay(实际上是线程池的延迟执行,后续涉及线程池再详细说明),传入了 $continuation

    • 判断逻辑 DelayKt.delay(1000L, (Continuation)$continuation) == var5 判断是否会被挂起,如果是则直接 return;

  2. 挂起点恢复时(线程池执行),调用 $continuation.resumeWith(内部调用了 invokSuspend,赋值挂起函数的执行结果),再次调用 mockNetworkData 方法;

  3. 挂起点恢复时再次调用 mockNetworkData,跳过匿名类对象,取出挂起点的执行结果(delay 返回结果是 Unit)和状态 label ,进入 Q2 处执行:

    • 恢复挂起点前的状态,也就是恢复 tmp 的值;

    • break 以后直接到 return 返回对应结果;

真正的挂起

难道 suspend 修饰的函数不是会被挂起吗?

并不是, suspend 仅仅表示函数 可以被挂起,但是否会被挂起还需要根据内部是否有挂起点,比如 delaylaunch 等内置的挂起函数;

来看下面例子:

  • E1 函数虽说被 suspend 修饰,但是并不会被挂起,idea 编译器也会提醒这个 suspend 关键字是多余的;

  • 反编译后的结果,直接返回 “123”,实际上这是 同步 执行的流程,因此不会被挂起也无需被挂起;

 // E1: 不会被挂起的挂起函数 
 suspend fun nonSuspend(): String {
     return "123"
 }
 ​
 // E1 反编译的结果
 public static final Object nonSuspend(@NotNull Continuation $completion) {
     return "123";
 }

为了判断是 同步 还是 异步(被挂起),通过函数的返回值是否为 getCOROUTINE_SUSPENDED() 来判断是否会被真正的挂起;