Kotlin协程

1 协程的基本用法

1.1 引入协程

module下的gradle文件中引入

implementation "org.jetbrains.kotlinx:kotlinx-coroutines-android:1.3.2"

1.2 launch用法

1.2.1 基本用法

override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        setContentView(R.layout.activity_main)
        logcat("main activity launch")
        GlobalScope.launch(Dispatchers.Main) {
            logcat("this is main coroutins")
        }
    }

打印

10-30 10:00:32.972 21012 21012 E xclog   : main activity launch
10-30 10:00:33.221 21012 21012 E xclog   : this is main coroutins

使用Dispatchers.XX 可以指定协程部分运行的线程, 如果修改为Dispatchers.IO, 则运营结果如下:

0-30 10:01:59.671 21952 21952 E xclog   : main activity launch
10-30 10:01:59.695 21952 22040 E xclog   : this is main coroutins

协程也可以delay, delay不会阻塞线程的运行

1.2.2 join的用法

利用Join可以等待一个协程执行完成

override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        setContentView(R.layout.activity_main)
        logcat("main activity launch")
        GlobalScope.launch(Dispatchers.Main) {
            logcat("main cor launch start")
            val job = GlobalScope.launch {
                logcat("this is coroutine scope")
            }
            job.join()
            logcat("main cor launche finished")
        }
    }

输出如下:

10-30 10:13:39.622 22815 22815 E xclog   : main activity launch
10-30 10:13:39.771 22815 22815 E xclog   : main cor launch start
10-30 10:13:39.774 22815 22929 E xclog   : this is coroutine scope
10-30 10:13:39.776 22815 22815 E xclog   : main cor launche finished

1.3 async用法

1.3.1 基本用法

与launch不同, async可以利用await可以返回一个协程的执行结果
代码

override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        setContentView(R.layout.activity_main)
        logcat("main activity launch")
        GlobalScope.launch(Dispatchers.Main) {
            logcat("main cor launch start")
            val deferred = GlobalScope.async {
                logcat("this is coroutine scope")
                45
            }
            logcat("main cor launche finished:${deferred.await()}")
        }
    }

执行结果:

10-30 10:16:10.922 23727 23727 E xclog   : main activity launch
10-30 10:16:11.061 23727 23727 E xclog   : main cor launch start
10-30 10:16:11.064 23727 23763 E xclog   : this is coroutine scope
10-30 10:16:11.066 23727 23727 E xclog   : main cor launche finished:45

1.4 suspend的用法

suspend用于声明一个函数\方法, 标明这个函数是只能在协程或者suspend函数中调用. suspend声明的函数体内, 可以试用delay, async, join等协程相关方法.
这里注意一点: delay, async, join等协程相关方法只能在协程或者suspend函数中调用
代码:

override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        setContentView(R.layout.activity_main)
        logcat("main activity launch")
        GlobalScope.launch(Dispatchers.Main) {
            logcat("this is cor:${getResult()}")
        }
    }

    private suspend fun getResult(): Int {
        delay(1000)
        return 100
    }

运行结果:

10-30 10:28:50.796 25272 25272 E xclog   : main activity launch
10-30 10:28:51.929 25272 25272 E xclog   : this is cor:100

2 协助用来解决回调地狱问题

2.1 回调地狱问题

假设一个题目:

  1. 在主线程中启动计算A, 计算A要持续10S的时间, 计算A不能阻塞主线程的运行
  2. 计算A完成后, 在主线程中显示计算A的结果
  3. 显示A结果完成后, 开始计算B, 计算B要持续4S的时间, 计算B不能阻塞主线程的运行
  4. 计算B完成后, 要在主线程中显示计算B的结果

2.1.1 传统的做法

private fun startTest() {
        // 在主线程中开启一个子线程开始计算
        logcat("start A cal")
        Thread {
            // 用sleep模拟计算
            Thread.sleep(10 * 1000)
            //计算完成后, 使用
            mHandler.post {
                logcat("start A cal finish")
                //开启计算B
                logcat("start B cal")
                Thread {
                    Thread.sleep(4 * 1000)
                    mHandler.post {
                        logcat("B cal finish")
                    }
                }.start()
            }
        }.start()
    }

这种做法, 一方面代码很难看, 另外一方面, 子线程启动后, 很难去进行控制
幸好我们有一个非常优秀的工具, Rxjava

2.1.2 Rxjava做法

通过下面代码可以看到, Rxjava已经优秀很多了, 整个调用特别有条理

private fun startByRxjava(){
        val sub = Observable.just("")
            .observeOn(Schedulers.newThread())
            .map {
                logcat("start A cal")
                Thread.sleep(10 * 1000)
                it
            }
            .observeOn(AndroidSchedulers.mainThread())
            .map {
                logcat("A cal finished")
                it
            }
            .observeOn(Schedulers.newThread())
            .map {
                logcat("start B cal")
                Thread.sleep(4 * 1000)
                it
            }
            .observeOn(AndroidSchedulers.mainThread())
            .subscribe {
                logcat("B cal finished")
            }
    }

2.1.3 协程的做法

做法1, delay和Sleep方法不同, 不会阻塞线程运行

private fun startV1(){
        GlobalScope.launch(Dispatchers.Main) {
            logcat("start cal A")
            delay(10 * 1000)
            logcat("cal A test finish")
            logcat("start cal B")
            delay(5 * 1000)
            logcat("cal B result")
        }
    }

上述做法计算都放在了主线程, 还是会"占用"的资源. 换一种做法

private fun startV2() {
        GlobalScope.launch(Dispatchers.Main) {
            startTestA()
            logcat("cal A finish")
            startTestB()
            logcat("cal B finish")
        }
    }

    private suspend fun startTestA(): Int {
        return GlobalScope.async {
            logcat("start cal A")
            kotlinx.coroutines.delay(10 * 1000)
            10
        }.await()
    }

    private suspend fun startTestB(): Int {
        return GlobalScope.async {
            logcat("start cal B")
            kotlinx.coroutines.delay(5 * 1000)
            10
        }.await()
    }

2.2 利用协程+retrofit2改造网络请求

我们经常使用RxJava+retrofit2来进行网络请求, 这里不赘述用法, 下面介绍下使用协程+retrofit2来改造网络请求

  1. gradle: 增加下面代码
api 'com.jakewharton.retrofit:retrofit2-kotlin-coroutines-adapter:0.9.2'
  1. Retrofit初始化的时候增加CoroutineCallAdapterFactory
val retrofit = Retrofit.Builder()
                .baseUrl(baseUrl)
                .addConverterFactory(GsonConverterFactory.create())
                .addCallAdapterFactory(RxJavaCallAdapterFactory.create())
                .addCallAdapterFactory(CoroutineCallAdapterFactory())
                .client(okHttpHttp).build()
        service = retrofit.create(ApiService::class.java)
  1. ApiSevice:
@GET("http://quan.lukou.com/api/taolijin/fetch")
    fun taoLiJinTaskV2(): Deferred<DataWrapper<TaoLiJinTaskBean>>
  1. 调用:
@Test
    fun testDefered() {
        GlobalScope.launch(context = Dispatchers.Main) {
            val result = ApiFactory.taoLiJiJinTaskV2().await()
            logcat("result:"+GsonManager.instance.toJson(result))
        }
    }

3 自定义协程上下文

在Dispatchers中, 定义了四种预置的上下文,来指定协程在哪个线程中执行,分别是

  • Default: 默认, 最大的线程数==cpu的核心数, 挂起后可能会在任何一个线程内恢复
  • Main: 主线程
  • Unconfined: 不限定线程. 协程从开始到第一个挂起位置部分会在调用线程内自行, 挂起恢复之后, 会在任意一个线程中恢复, 如下:
GlobalScope.launch(Dispatchers.Unconfined){
            logcat("abc")
            kotlinx.coroutines.delay(1099)
            logcat("bcd")
            kotlinx.coroutines.delay(200)
            logcat("cdf")
        }

执行结果如下:

10-30 14:18:05.725 13635 13635 E xclog   : abc
10-30 14:18:06.827 13635 13669 E xclog   : bcd
10-30 14:18:07.027 13635 13669 E xclog   : cdf
  • IO: 最大线程数为64个, 挂起后会在任何一个线程内恢复

我们仿照系统定义, 自定义一个协程调度器

object SingleThreadDispather : ExecutorCoroutineDispatcher() {
    private var myExecutor: Executor? = null
    override val executor: Executor
        get() = run {
            if (myExecutor == null) {
                myExecutor = Executors.newSingleThreadExecutor()
            }
            myExecutor!!
        }

    override fun close() {
    }

    override fun dispatch(context: CoroutineContext, block: Runnable) {
        executor.execute(block)
    }
}

4 Flow的用法

基本用法

我们如果想让一个函数返回多个值, 应该怎么办? 很简单, 用List或者Array

override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        setContentView(R.layout.activity_main)
        logcat("main activity launch")
        getResult().forEach {
            logcat("foreach:${it}")
        }
    }

    private fun getResult(): List<Int> {
        val arrayList = ArrayList<Int>()
        for (i in 0..3) {
            Thread.sleep(100) // 假设每一个值的计算都需要花费100毫秒
            arrayList.add(i)
        }
        return arrayList
    }

这种做法的坏处是显而易见的

换一种写法

override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        setContentView(R.layout.activity_main)
        logcat("main activity launch")
        getResult().forEach {
            logcat("foreach:${it}")
        }
    }
    private fun getResult1(): Sequence<Int> = sequence {
        for (i in 0..3) {
            Thread.sleep(100)
            logcat("yield:${i}")
            yield(i)
        }
    }

输出为

11-14 20:33:25.477 32187 32187 E xclog   : main activity launch
11-14 20:33:25.580 32187 32187 E xclog   : yield:0
11-14 20:33:25.581 32187 32187 E xclog   : foreach:0
11-14 20:33:25.681 32187 32187 E xclog   : yield:1
11-14 20:33:25.681 32187 32187 E xclog   : foreach:1
11-14 20:33:25.781 32187 32187 E xclog   : yield:2
11-14 20:33:25.781 32187 32187 E xclog   : foreach:2
11-14 20:33:25.881 32187 32187 E xclog   : yield:3
11-14 20:33:25.882 32187 32187 E xclog   : foreach:3

已经好多了, 是吧, 但是这里只能用sleep, 还是会阻塞线程

如果用Flow

override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        setContentView(R.layout.activity_main)
        logcat("main activity launch")
        GlobalScope.launch(Dispatchers.Main) {
            getResult2().collect { logcat("collect $it") }
        }
    }

    private fun getResult2():Flow<Int> = flow {
        for(i in 0..3){
            delay(100)
            logcat("emit $i")
            emit(i)
        }
    }

输出为:

11-14 20:38:32.796 32513 32513 E xclog   : emit 0
11-14 20:38:32.796 32513 32513 E xclog   : collect 0
11-14 20:38:32.897 32513 32513 E xclog   : emit 1
11-14 20:38:32.897 32513 32513 E xclog   : collect 1
11-14 20:38:32.998 32513 32513 E xclog   : emit 2
11-14 20:38:32.998 32513 32513 E xclog   : collect 2
11-14 20:38:33.098 32513 32513 E xclog   : emit 3
11-14 20:38:33.098 32513 32513 E xclog   : collect 3

Flow的其他用法

  1. 创建flow

除了上面创建flow外, 还有其他办法

GlobalScope.launch(Dispatchers.Main) {
            (0..3).asFlow().collect { logcat("collect1 $it") }
            listOf(1, 2, 4).asFlow().collect { logcat("collect 2: $it") }
            flowOf(4, 5, 6).collect { logcat("collect 3:$it") }
        }
  1. map

对flow中的每个值做处理

(4..7).asFlow()
                .map {
                    (it + 1).toString()
                }
                .collect {
                    logcat("collect 4:"+it)
                }
  1. zip

将两个flow组合

(4..8).asFlow()
                .map {
                    (it + 1).toString()
                }
                .zip((8..9).asFlow()) { a, b ->
                    "a:${a} b:${b}"
                }
                .collect {
                    logcat("collect 4:" + it)
                }
  1. reduce

实现累加的效果

val sum = (4..8).asFlow()
                .reduce { accumulator, value ->  accumulator + value}
            logcat("sum :$sum")
  1. flatMapConcat

可以将每个值作为参数传递到下个flow中, 生成一个新的flow

GlobalScope.launch(Dispatchers.Main) {
            (0..3).asFlow().collect { logcat("collect1 $it") }
            listOf(1, 2, 4).asFlow().collect { logcat("collect 2: $it") }
            flowOf(4, 5, 6).collect { logcat("collect 3:$it") }
            val sum = (1..2).asFlow()
                .flatMapConcat {
                    getResult3(it)
                }
                .reduce { accumulator, value -> accumulator + value }
            logcat("sum :$sum")
        }
	private fun getResult3(value: Int): Flow<Int> = flow {
        emit(value + 1)
    }

5 Channel的用法

利用Channel可以实现协程间的通信

@Test
    fun testChannel() {
        logcat("testChannel")
        GlobalScope.launch {
            repeat(5) {
                logcat("channel:${mChannel.receive()}")
            }
        }
        GlobalScope.launch {
            for (i in 1..5) {
                logcat("send $i")
                mChannel.send(i)
            }
        }
        Thread.sleep(5000)
    }

Channel适合来做EventBus, 来实现一个Channel版本的EventBus

Logo

开源鸿蒙跨平台开发社区汇聚开发者与厂商,共建“一次开发,多端部署”的开源生态,致力于降低跨端开发门槛,推动万物智联创新。

更多推荐