
一、概念协程必须运行在一个线程上所以要指定调度器。是一个抽象类Dispatcher是一个标准库中帮我们封装了切换线程的帮助类可以调度协程在哪类线程上执行。创建协程时上下文如果没有指定也没有继承到调度器则会添加一个默认调度器调度器通过 ContinuationInterceptor 延续体拦截器实现的。Dispatcher 只是指定作用域内代码执行在哪类线程上默认情况下内部是单线程执行如指定IO挂起前在A线程恢复后在B线程本质还是单线程只有在多个协程同时调用该代码才存在多线程并发问题。二、模式由于子协程会继承父协程的上下文在父协程上指定调度器模式后子协程默认使用这个模式。DEFAULT 针对的是 CPU 密集型计算任务CPU 核心数是就是最大线程数量满负荷运行再增加线程数反而降低效率IO 针对的是发生在内存之外的磁盘网卡在做输入输出此时 CPU 是空闲的更大的线程数量能充分利用 CPU 调度。IO 和 DEFAULT 模式共享同一线程池重用线程起到优化DEFAULT 切换到 IO 大概率停留在同一线程上两者对线程数量限制是独立的不会让对方饥饿。最大限度一起使用的话默认同时活跃的线程数为 64CPU数。如果切换的模式和当前模式是相同的是不会再切线程而是保持在当前线程继续执行后续代码。但 Main 由于底层采用的是 Handler.post()因此提供了 immediate。Dispatcher.Main运行于主线程在 Android 中就是 UI 线程用来处理一些 UI 交互的轻量级任务。更新UIDispatcher.Main.immediate当我们已经处在主线程时再调度到主线程是不必要的开销指定为 immediate 只会在非主线程时调度否则直接执行。(LifecycleScope、ViewModelScope 就处在Android默认的主线程中因此上下文中的调度器使用了这个)更新UIDispatcher.IO运行于线程池专为IO阻塞型等待型任务进行了优化。最大线程数为64个只要没超过且没有空闲线程就一直可开辟新线程执行新任务。数据库文件读写网络处理Dispatcher.Default运行于线程池专为CPU密集型计算任务进行了优化。最大线程数为CPU核心数但不少于2个若全在忙碌时新任务无法得到执行。数组排序Json解析处理差异判断计算BitmapDispatcher.Unconfined不改变线程在启动它的线程执行在恢复它的线程执行也就是当前线程在哪就在哪执行。调度成本最低性能最好但存在不可控风险如处在主线程执行了阻塞操作实际开发不会用到。当不需要关心协程在哪个线程上被挂起时使用。三、自定义线程数3.1 限制线程数 limitedParallelism()1.6版本引入。对于Default模式当有一个开销很大的任务可能会导致其它使用相同调度器的协程抢不到线程执行权这个时候就可以用来限制该协程的线程使用数量。对于IO模式当有一个开销很大的任务可能会导致阻塞太多线程让其它任务暂停等待突破默认64个线程的限制加速执行不显著。传参将线程限制为1解决多线程并发修改数据的同步问题。但如果阻塞了它其它操作都要等待。public open fun limitedParallelism(parallelism: Int): CoroutineDispatchersuspend fun main(): Unit coroutineScope { //使用默认IO模式 launch { printTime(Dispatchers.IO) //打印Dispatchers.IO 花费了: 2038/ } //使用limitedParallelism增加线程 launch { val dispatcher Dispatchers.IO.limitedParallelism(100) printTime(dispatcher) //打印LimitedDispatcher1cc12797 花费了: 1037 } } suspend fun printTime(dispatcher: CoroutineDispatcher) { val time measureTimeMillis { coroutineScope { repeat(100) { launch(dispatcher) { Thread.sleep(1000) } } } } println($dispatcher 花费了: $time) }3.2 专用的单线程/线程池不推荐在没有 limitedParallelism() 以前就是这样做的。专用的线程可能会抵消地使用未使用的线程保持活跃状态却不与其它业务共享这些线程使用完容易忘记关闭创建太多消耗系统资源容易用完忘记关闭造成内存泄漏。3.2.1 单线程 newSingleThreadContext是协程提供的一个用于创建单线程调度器的函数它可以确保所有在该调度器上执行的协程都在同一个线程中顺序执行。public fun newSingleThreadContext(name: String): CloseableCoroutineDispatcher newFixedThreadPoolContext(1, name)val singleThreadContext newSingleThreadContext(SingleThreadContext) //两个launch不再是并行执行而是按顺序 launch(singleThreadContext) {...} launch(singleThreadContext) {...} singleThreadContext.close()3.2.2 线程池 newFixedThreadPoolContext是协程提供的一个用于创建固定大小线程池调度器的函数当所有线程都忙时新任务会排队等待。public expect fun newFixedThreadPoolContext(nThreads: Int, name: String): CloseableCoroutineDispatcherval threadPoolContext newFixedThreadPoolContext(3, ThreadPoolContext) launch(threadPoolContext) {...} threadPoolContext.close()3.3.3 Java线程转换 asCoroutineDispatcher()将 Java 的方式转为协程版本。val singleThreadExecutor Executors.newSingleThreadExecutor().asCoroutineDispatcher() val fixedThreadPool Executors.newFixedThreadPool(3).asCoroutineDispatcher() launch(singleThreadExecutor) {}五、多线程并发问题某个协程对共享变量的更新可能不会立即被其他协程看见尤其是在多线程环境下协程可能运行在不同的线程上。多个协程可能会交错执行这些步骤导致最终值错误。协程的挂起与恢复可能导致锁资源竞争。如果协程挂起期间未正确释放锁可能造成死锁。这与我们在使用线程时是一样的创建10个协程每个协程执行 i 1000次预期结果 i10000实际会小于这个值。i 不是线程安全操作它包含了读取、计算、写入。//不指定Dispatchers的话默认是BlockingEventLoop不是Default fun main(): Unit runBlocking(Dispatchers.Default) { var count 0 val time measureTimeMillis { repeat(10) { //创建10个协程 launch { repeat(100) { //每个协程执行运算100次 count } } } } print(count $count, time $time) //打印count800time6 }4.1 避免使用共享变量fun main(): Unit runBlocking { var count 0 val deferreds mutableListOfDeferredInt() val time measureTimeMillis { repeat(10) { val deferred async(Dispatchers.Default) { var i 0 repeat(1000) { i } returnasync i } deferreds.add(deferred) } deferreds.forEach { count it.await() } } print(i $count, time $time) } //打印i 10000耗时774.2 使用Java方法不推荐可以使用 Synchronized、Lock、Atomic。由于是线程模型下的阻塞方式不支持调用挂起函数会影响协程挂起特性。4.2.1 使用同步锁fun main() runBlocking { val start System.currentTimeMillis() var i 0 val jobs mutableListOfJob() Synchronized fun add() { i } repeat(10) { val job launch(Dispatchers.Default) { repeat(1000) { add() } } jobs.add(job) } jobs.joinAll() println(i $i耗时${System.currentTimeMillis() - start}) } //打印i 10000耗时714.2.2 使用同步代码块fun main() runBlocking { val start System.currentTimeMillis() val lock Any() var i 0 val jobs mutableListOfJob() repeat(10) { val job launch(Dispatchers.Default) { repeat(1000) { synchronized(lock) { i } } } jobs.add(job) } jobs.joinAll() println(i $i耗时${System.currentTimeMillis() - start}) } //打印i 10000耗时734.2.3 使用可重入锁 ReentrantLockfun main() runBlocking { val start System.currentTimeMillis() val lock ReentrantLock() var i 0 val jobs mutableListOfJob() repeat(10) { val job launch(Dispatchers.Default) { repeat(1000) { lock.lock() i lock.unlock() } } jobs.add(job) } jobs.joinAll() println(i $i耗时${System.currentTimeMillis() - start}) } //打印i 10000耗时834.2.4 使用 Atomic 保证原子性fun main() runBlocking { val start System.currentTimeMillis() var i AtomicInteger(0) val jobs mutableListOfJob() repeat(10) { val job launch(Dispatchers.Default) { repeat(1000) { i.incrementAndGet() } } jobs.add(job) } jobs.joinAll() println(i $i耗时${System.currentTimeMillis() - start}) } //打印i 10000耗时894.3 使用单线程不推荐fun main() runBlocking { val start System.currentTimeMillis() val mySingleDispatcher Executors.newSingleThreadExecutor { Thread(it, 我的线程).apply { isDaemon true } }.asCoroutineDispatcher() var i 0 val jobs mutableListOfJob() repeat(10) { val job launch(mySingleDispatcher) { repeat(1000) { i } } jobs.add(job) } jobs.joinAll() println(i $i耗时${System.currentTimeMillis() - start}) mySingleDispatcher.close() //用完容易忘记关闭 } //打印i 10000耗时644.4 使用 MutexJava方式不支持调用挂起函数同步锁是阻塞式的会影响协程特性为此 Kotlin 提供了非阻塞式互斥锁Mutex会挂起协程让线程自由执行其他工作。使用 mutex.lock() 和 mutex.unlock() 包裹需要同步的计算逻辑就可以实现多线程同步了但由于包裹内容可能出现的异常使得 unlock() 无法被执行写在 finally{} 中会很繁琐因此提供了扩展函数 mutex.withLock{ }本质就是在 finally{ } 中调用了 unlock()确保同一时间只有一个协程可以执行该区域。public suspend inline fun T Mutex.withLock(owner: Any? null, action: () - T): T {lock(owner)try {return action()} finally { // 注意这里并没有 catch 代码块所以不会捕获异常unlock(owner)}}fun main() runBlocking { val start System.currentTimeMillis() var i 0 val mutex Mutex() //使用方式一 mutex.lock() // try { // repeat(10000) { i } // } catch (e: Exception) { // e.printStackTrace() // } finally { // mutex.unlock() // } //使用方式二 mutex.withLock { try { repeat(10000) { i } } catch (e: Exception) { e.printStackTrace() } } println(i $i耗时${System.currentTimeMillis() - start}) } //方式一打印i 10000耗时17 //方式二打印i 10000耗时174.5 使用 Channel 的 actor()Channel 默认配置是并发安全的生产和消费是挂起函数会同步交替进行。actor() 创建 SendChannel 定义消费逻辑往流中发送一次消息就执行一次。sealed class Msg { object AddMsg : Msg() class ResultMsg(val result: CompletableDeferredInt) : Msg() } OptIn(ObsoleteCoroutinesApi::class) fun main() runBlocking { val start System.currentTimeMillis() val actor actorMsg { var i 0 for (msg in channel) { when (msg) { is Msg.AddMsg - i is Msg.ResultMsg - msg.result.complete(i) } } } val jobs mutableListOfJob() repeat(10) { val job launch { repeat(1000) { actor.send(Msg.AddMsg) } } jobs.add(job) } jobs.joinAll() val deferred CompletableDeferredInt() actor.send(Msg.ResultMsg(deferred)) val result deferred.await() actor.close() println(i $result耗时${System.currentTimeMillis() - start}) } //打印i 10000耗时1674.6 使用 Semaphore信号量用于限制并发访问某个资源的协程数量避免资源耗尽。如在 Android 开发中你可能有 100 个网络请求任务要执行但希望最多只有 8 个请求同时进行以减轻服务器和客户端的压力。一个 Semaphore 的许可数为 1 时它的行为就等同于一个 Mutex互斥锁。Mutex 常用于保护资源的独占访问、Semaphore 常用于限制并发数量、Channel 常用语传递数据。fun main() runBlocking { var count 0 val semaphore Semaphore(1) //只允许一个协程访问 repeat(1000) { GlobalScope.launch { semaphore.withPermit { count } } }.joinAll() println(count) }五、第三方库调用的线程问题使用第三方库的时候Retrofit网络请求、Room查询数据库它们已经在内部处理好了IO调度用它们自己的线程来管理因此不要用 withContext(Dispatcher.IO) 去包裹你只是把诸如准备请求和解析 json 这类占用 cpu 的工作交给了 IO 调度器反而画蛇添足。IO任务是一种不会给CPU带来高负载的工作大部分时间都消耗在等待硬件响应上持久化存储或远程主机。CPU密集型任务主要靠CPU运算默认的 Dispatcher.Default 线程数和CPU核心数相对应能充分利用CPU核心避免过多的线程频繁切换上下文增加额外开销本来一个线程就能高效完成结果 Dispatcher.IO 的64个线程抢着做。