位置:首页 > Kotlin > 深入理解Kotlin协程select多路复用与并发安全策略

深入理解Kotlin协程select多路复用与并发安全策略

时间:2026-08-14  |  作者:实验室老王  |  阅读:0

为什么需要 IO 多路复用?

我们先从一个简单的场景开始,来了解什么是阻塞 IO(Blocking IO):

深入理解Kotlin协程(十一):以逸待劳,探秘select多路复用与并发安全策略

当你调用 read() 从 Socket 获取数据时,进程会调用 read()

内核发现数据还未到来,它就会挂起阻塞这个进程。只有当数据到来时,内核才会唤醒进程,进程才能继续执行并拿到数据。

问题是,同一时间有 100 个客户端连接怎么办?

难道要开 100 个线程,让每个线程都阻塞等待吗?

所以我们就会想到:如果一个线程可以同时等待多个 IO 就好了。于是,就有了 IO 多路复用。

IO 多路复用到底是什么?

简单来说,就是一个线程同时监视多个文件描述符

只要任意一个就绪,就通知我们来处理。

这场景像不像某些预制菜外卖店,一个“大厨”同时盯着多个微波炉。

只要某台微波炉发出了“叮”的一声,就会过去打包出餐。

其中,多路指的是多个 IO 流或多个 Socket。

复用则是复用同一个线程/进程去处理,复用的是等待的能力。

操作系统的 select 机制

而 select 就是一个内核级支持的 IO 多路复用机制,它的工作逻辑是:

把当前需要监视的 Socket 集合交给内核之后,接下来调用 select 进入阻塞状态。

这时候进程会休眠,不再占用 CPU。

真正的轮询工作由内核来完成。它会持续检查所有 Socket。

只要其中有一个变成就绪状态——无论是可读还是可写——select 就会立刻返回。

然后再去遍历这些 Socket,找出那个已经就绪的对象并进行处理。

它有很多与生俱来的缺陷,因此后续还出现了 poll,出现了 epoll。

协程中的多路复用

这里先说明一点:操作系统的 select 是内核级监控文件描述符。

而我们接下来要讲的 Kotlin 协程的 select 则是用户态的机制。

它监控的是协程间的通信事件,复用的是协程的挂起和恢复能力,而不是线程的阻塞。

回到本节的主题,我们来看看多路复用在协程中到底怎么用。

Deferred 的 await 复用

首先是 await 可以进行复用。

想象这样一个场景:当应用数据可以从本地或网络获取时,我们只希望展示更先返回的数据。

这时,可以使用 select 这样做:

fun main(): Unit = runBlocking {
val userId = "123456"
repeat(2) {
val localDeferred = getUserFromLocal(userId)
val remoteDeferred = getUserFromRemote(userId)

val userResult = select {
localDeferred.onAwait { it }
remoteDeferred.onAwait { it }
}

userResult.let {
println("user result: $userResult")
} : run {
// 获取网络数据展示最终结果
val remoteUseResult = remoteDeferred.await()
println("user result: $remoteUseResult")
}

// 拿到结果后,主动取消所有任务,防止协程泄漏
localDeferred.cancel()
remoteDeferred.cancel()
}
}

// 模拟本地缓存
object UserCache {
private val cache = ConcurrentHashMap()
fun get(userId: String): String {
return cache[userId]
}

fun put(userId: String, value: String) {
cache[userId] = value
}
}

fun CoroutineScope.getUserFromLocal(userId: String): Deferred = async(Dispatchers.IO) {
delay(Random.nextLong(1000, 2000).milliseconds)
UserCache.get(userId)
}

fun CoroutineScope.getUserFromRemote(userId: String): Deferred = async(Dispatchers.IO) {
delay(Random.nextLong(1000, 3000).milliseconds)
val userData = "User-$userId's data(remote)"
UserCache.put(userId, "User-$userId's data(local)")
userData
}

我们调用了 DeferredonAwait 函数。

select 中注册回调后,select 会调用最先返回的事件。

Channel 的读取复用

复用 Channel 和复用 await 类似:

fun main(): Unit = runBlocking {
val channels = List(10) { Channel() }
// 随机发送数据
launch {
delay(300.milliseconds)
channels[Random.nextInt(0, channels.size)].send(1)
}

// 接收最先收到的数据
val result = select {
channels.forEachIndexed { index, channel ->
// 或者调用 onReceive,它会将异常直接抛出
channel.onReceiveCatching { result ->
result.onSuccess { data ->
return@onReceiveCatching "Channel [${index + 1}] ==> $data"
}
// 当通道被关闭或是出现异常时返回 null 作为兜底,避免程序直接崩溃
return@onReceiveCatching null
}
}
}
println(result)
}

SelectClauseN 接口解析

能被 select 的事件,都实现了 SelectClauseN 接口。

例如刚刚的 onReceiveCatching 成员属性类型是 SelectClause1>

事件有以下几种类型:

  • SelectClause0: 事件没有返回值,例如 Job.join() 没有返回值,因此其 onJoin 成员属性是 SelectClause0 类型。
  • SelectClause1: 事件有返回值,例如上面的 onAwaitonReceive
  • SelectClause2: 事件有返回值,同时还需额外的参数。

以 Channel 的 onSend 为例,它有一个 param 参数和一个 block 返回值。

当指定的参数 param 成功被发送到通道时,block 回调将会触发。

block 中的参数是成功发送到的 Channel 对象。

官方示例:

fun main(): Unit = runBlocking {
val sendChannels = List(4) { index ->
Channel(
onUndeliveredElement = {
println("Undelivered element $it for $index")
}
).also { channel ->
launch {
withTimeout(1.seconds) {
println("Consumer $index receives: ${channel.receive()}")
}
}
}
}
val element = 42
select {
for (channel in sendChannels) {
channel.onSend(element) {
println("Sent to channel $it")
}
}
}
}

上述代码会随机消费一个 Channel。

此时对应的 Channel 就会成功发送数据,并触发成功回调。

剩余的三个 Channel 会被 select 取消,走到 onUndeliveredElement 回调。

使用 Flow 实现复用

使用 Flow 也可以实现上述多路复用的效果。

核心思路是:只取第一个任务,收到第一个任务后,手动取消其余 Job。

fun main(): Unit = runBlocking {
val userId = "123456"
// 并发调用获取 Deferred 列表
val deferredList = listOf(
getUserFromLocal(userId),
getUserFromRemote(userId)
)

try {
deferredList.map { deferred ->
// 创建单独的 Flow 来获取结果
flow {
val result = deferred.await()
emit(result)
}
}
.merge()
.first().also {
println(it)
}
} finally {
// 取消剩余任务
deferredList.forEach {
it.cancel()
}
}
}

同理,Channel 的读取复用也可以这样处理:

fun main(): Unit = runBlocking {
val size = 10
val channels = List(size) {
Channel()
}
val flows = channels.map {
it.consumeAsFlow()
}.merge()

launch {
val index = Random.nextInt(size)
channels[index].send("67 from Channel [$index]")
}

val result = flows.first()
println(result)
}

协程的并发安全

为什么会有并发安全问题?

Kotlin 协程也存在着并发安全,因其运行在 Ja va 平台上,最终会被调度到某个线程上执行。

因此,count++ 是不安全的。我们来看一个简单的计数问题:

fun main(): Unit = runBlocking {
var count = 0
// 让这个数字尽可能大
List(100000) {
launch(Dispatchers.Default) {
count++
}
}.joinAll()
println("the final count: $count")
}

结果可能是:99849。

原因我们都很清楚了:

  • count 变量不可见,它的读写不会立即同步到主内存。
  • count++ 不是原子操作,读、改、写的过程中会被其他线程插入。

如果不清楚的话,可以看我的这篇博客:Ja va 多线程指南:从基础用法到线程安全

协程并发安全的解决方案

为了解决这个问题,我们可以将 count 声明为原子类,也可以直接加锁。

但在协程中,我们有更好的解决方案:

  • Channel: 并发安全的消息通道。
  • Mutex: 轻量级锁,语义上与线程锁类似,但获取不到锁时,它不会阻塞线程,只是会挂起等待锁释放。
fun main(): Unit = runBlocking {
var count = 0
val mutex = Mutex()
List(100000) {
launch {
// 封装了 lock 和 unlock 操作
mutex.withLock {
count++
}
}
}.joinAll()
println("the final count: $count")
}
  • Semaphore: 轻量级信号量,控制同一时间访问资源的协程最大并发数量。当参数为 1 时,效果等价于使用 Mutex。
fun main(): Unit = runBlocking {
// 最多允许两个协程同时访问
val semaphore = Semaphore(2)
val jobs = List(5) { taskId ->
launch(Dispatchers.Default) {
// 获取许可证
semaphore.acquire()
try {
println("Starting task $taskId, thread: ${Thread.currentThread().name}")
delay(3000.milliseconds)
println("End task $taskId, thread: ${Thread.currentThread().name}")
} finally {
// 归还许可证
semaphore.release()
}
}
}
jobs.joinAll()
}

上面这段代码,展示的就是 Semaphore 最原始的获取与释放机制。

到了实际开发里,通常更推荐直接用 Semaphore.withPermit { ... } 这个扩展函数。

它会把许可证的释放自动处理掉,省心不少。

它的用法也很直观,和 MutexwithLock 基本完全一样。

总结:最好的并发就是不共享状态

实际上,大多时候我们都无需硬刚线程安全问题。

通过让协程访问的共享资源不可变,让访问外部状态的函数变为纯函数,这样就能避免并发安全。

以前面的计数为例:

fun main(): Unit = runBlocking {
val count = 0
val sum = List(100000) {
async {
1
}
}.awaitAll().sum()
val result = count + sum
println("the result is $result")
}

运行结果:"the result is 100000"

免责声明:文中图文均来自网络,如有侵权请联系删除,心愿游戏发布此文仅为传递信息,不代表心愿游戏认同其观点或证实其描述。

相关文章

更多

精选合集

更多

大家都在玩

热门话题

大家都在看

更多