Kotlin协程
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 回调地狱问题
假设一个题目:
- 在主线程中启动计算A, 计算A要持续10S的时间, 计算A不能阻塞主线程的运行
- 计算A完成后, 在主线程中显示计算A的结果
- 显示A结果完成后, 开始计算B, 计算B要持续4S的时间, 计算B不能阻塞主线程的运行
- 计算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来改造网络请求
- gradle: 增加下面代码
api 'com.jakewharton.retrofit:retrofit2-kotlin-coroutines-adapter:0.9.2'
- 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)
- ApiSevice:
@GET("http://quan.lukou.com/api/taolijin/fetch")
fun taoLiJinTaskV2(): Deferred<DataWrapper<TaoLiJinTaskBean>>
- 调用:
@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的其他用法
- 创建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") }
}
- map
对flow中的每个值做处理
(4..7).asFlow()
.map {
(it + 1).toString()
}
.collect {
logcat("collect 4:"+it)
}
- zip
将两个flow组合
(4..8).asFlow()
.map {
(it + 1).toString()
}
.zip((8..9).asFlow()) { a, b ->
"a:${a} b:${b}"
}
.collect {
logcat("collect 4:" + it)
}
- reduce
实现累加的效果
val sum = (4..8).asFlow()
.reduce { accumulator, value -> accumulator + value}
logcat("sum :$sum")
- 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
更多推荐



所有评论(0)