协程的协作等待:替代 CountDownLatch
在 Java 线程中,当一个线程【A】需要等待多个其他线程【B】的完成才能继续时,我们会使用 CountDownLatch(Latch 的意思是门闩)。
使用方式:首先需要初始化它的计数值,在线程【A】中调用 await(),然后在每个线程【B】完成时调用 countdown() 即可。
只有当调用次数达到预设值时,await() 才会返回。这就实现了只有所有线程【B】执行完毕后,线程【A】才能接着执行的效果。
fun main() {
val count = 2
val latch = CountDownLatch(count)
val b1 = thread {
Thread.sleep(1000)
latch.countDown()
}
val b2 = thread {
Thread.sleep(3000)
latch.countDown()
}
val a = thread {
latch.await()
println("Thread A finished")
}
// 等待线程 A 完成
a.join()
}
Job.join() 与 Deferred.await()
在协程中,我们当然会想到 Job.join() 与 Deferred.await(),这两个函数都可用于等待协程的完成。
fun main(): Unit = runBlocking {
// 协程的等待
val job1 = launch {
delay(2000)
}
val job2 = launch {
delay(3000)
}
val job3 = launch {
// 等待 job1 和 job2 结束
job1.join()
job2.join()
delay(1000)
}
job3.join()
// 协程的等待
val deferred1 = async {
delay(500)
return@async "async1"
}
val deferred2 = async {
delay(3500)
return@async "async2"
}
launch {
// 获取 deferred1 和 deferred2 的结果
val result1 = deferred1.await()
val result2 = deferred2.await()
delay(300)
println("launch4 result: $result1 $result2")
}
}
Job.join() 的写法在线程中非常少见,因为线程的管理成本很高,且没有协程中的结构化取消。
Channel
在形式上,只关心完成的次数,Channel 与 CountDownLatch 更接近:
fun main(): Unit = runBlocking {
val count = 2
val channel = Channel<Unit>(capacity = count)
launch { // 协程A
repeat(count) {
channel.receive() // 挂起等待 count 次
}
delay(300)
println("A Done")
}
launch { // 协程B1
delay(1000)
println("B1 Done")
channel.send(Unit)
}
launch { // 协程B2
delay(2000)
println("B2 Done")
channel.send(Unit)
}
}
select 表达式:先到先得
对于多个任务,只想获得最快完成的结果时,就可以使用 select 表达式。
比如说我们可以给 Job 设置 onJoin 回调:
fun main(): Unit = runBlocking {
val job1 = launch {
delay(1000)
println("job1 done")
}
val job2 = launch {
delay(2000)
println("job2 done")
}
val result = select {
job1.onJoin {
"job1"
}
job2.onJoin {
"job2"
}
}
println("The first completed coroutine is $result")
}
共享变量与互斥锁:Mutex 与 synchronized
如果对 Java 多线程不了解,可以看我的这篇博客:Java 多线程指南:从基础用法到线程安全
竞争条件 (Race Condition),又称竞态条件,指的是多个线程访问或修改共享资源时,由于执行顺序的不确定性,导致程序的行为不一致或出现错误。在并发编程中,这是最常见且难以调试的 Bug 来源之一。当两个或多个线程同时读写同一块内存区域,且至少有一个线程在执行写操作时,如果没有适当的同步机制保护,最终的结果将依赖于线程调度的具体时序,从而产生不可预测的行为。
比如运行这段代码,发现 count 最终的值竟然不是零。这种现象通常发生在多个线程同时对一个共享计数器进行“读取-修改-写入”操作时。假设初始值为零,线程 A 读取到零,线程 B 也读取到零,随后 A 将值加一写回,B 也将值加一写回。理论上结果应为二,但由于缺乏原子性保护,实际结果可能因覆盖操作而变为零或一。这种错误并非每次运行都会出现,而是具有随机性,这使得问题排查变得极其困难。
为了解决这一问题,我们需要引入互斥锁(Mutex)和同步关键字(synchronized)。互斥锁是一种同步原语,它确保在同一时刻只有一个线程可以访问被保护的共享资源。当一个线程获取了锁之后,其他试图获取该锁的线程将被阻塞,直到锁被释放。在 Java 中,synchronized 关键字提供了内置的锁机制,它可以修饰方法或代码块,确保被修饰的代码段具有原子性。通过这种方式,我们可以强制多个线程串行化地访问共享变量,从而消除竞争条件,保证程序行为的一致性和正确性。
除了使用内置锁,我们还可以利用更高级的并发工具类,如 ReentrantLock 或 AtomicInteger,来实现更细粒度的控制。这些工具提供了非阻塞算法或更灵活的锁策略,能够在高并发场景下提供更好的性能。然而,无论使用何种工具,核心原则始终不变:确保对共享状态的访问是互斥的,或者通过不可变对象来避免共享状态。理解并正确应用这些同步机制,是编写健壮并发程序的关键所在。
在实际开发中,过度使用锁会导致性能下降,甚至引发死锁。因此,设计并发系统时,应尽可能减少共享状态的范围,优先使用线程局部变量或不可变对象。只有在确实需要共享可变状态时,才谨慎地使用同步机制。通过合理的设计,我们可以在保证线程安全的同时,最大化程序的并发性能。
协程中的 ThreadLocal 陷阱与解决方案
这时,我们可以使用 Java 的 synchronized 关键字(在 Kotlin 中是一个函数),这样临界区将会被多个线程所互斥。
fun main(): Unit = runBlocking {
var count = 0
val lock = Any()
val scope = CoroutineScope(Dispatchers.Default) // 使用多线程调度器
val job1 = scope.launch {
repeat(1_000_000) {
synchronized(lock = lock) {
count++
}
}
}
val job2 = scope.launch {
repeat(1_000_000) {
synchronized(lock = lock) {
count--
}
}
}
job1.join()
job2.join()
println("the final count is $count")
}
虽然可以在协程用 synchronized,因为 Kotlin 协程的本质是线程,最终协程的代码还是运行在线程中的。synchronized 会锁定当前的线程,自然也会卡住在该线程上运行的协程。
但并不推荐使用。因为你使用 synchronized,协程在等待锁时,会阻塞整个线程,而不是让出当前线程,导致了资源浪费。
使用 Mutex 互斥锁
在协程中,我们会使用 Mutex 互斥锁,相当于 Java 中的 Lock。
fun main(): Unit = runBlocking {
var count = 0
val mutex = Mutex()
val scope = CoroutineScope(Dispatchers.Default) // 使用多线程调度器
val job1 = scope.launch {
repeat(1_000_000) {
mutex.lock() // 挂起函数
try {
count++
} catch (e: Exception) {
e.printStackTrace()
} finally {
mutex.unlock() // 保证锁的释放
}
}
}
val job2 = scope.launch {
repeat(1_000_000) {
mutex.lock() // 挂起函数
try {
count--
} catch (e: Exception) {
e.printStackTrace()
} finally {
mutex.unlock() // 保证锁的释放
}
}
}
job1.join()
job2.join()
println("the final count is $count")
}
也可以使用便捷函数 withLock,它会自动上锁和解锁。
Mutex 与 synchronized 的区别在于:当锁被占用时,synchronized 会阻塞线程;而 mutex.lock() 会挂起协程,释放所占用的线程,让线程可以去干别的事。
锁的选择
实际中,该如何选择?
-
在纯协程环境中,永远使用
Mutex,因为不卡线程。 -
在纯线程环境中,就使用
synchronized或ReentrantLock。 -
如果共享变量既会被协程访问,又会被线程访问,就统一使用
synchronized或ReentrantLock,因为mutex.lock()是挂起函数,无法在线程中使用。
协程中的局部变量机制
在 Java 中,线程拥有独立的“局部变量”空间。然而,如果我们在协程中直接复用 ThreadLocal,可能会因为线程复用导致数据覆盖或丢失,引发难以排查的 Bug。
因此,不能简单地直接使用 ThreadLocal 来存储协程特定的状态。协程的局部变量机制依赖于 CoroutineContext,这是一种特殊的线程绑定方式。当协程与 ThreadLocal 进行交互时,推荐使用 asContextElement 来确保数据隔离。
实际上,asContextElement 的核心作用在于保证协程运行线程的自动同步。每当协程发生线程切换时,系统会自动保存并恢复 ThreadLocal 的值。这种机制确保了 ThreadLocal 在不同线程间迁移时,其关联的数据状态依然保持一致,从而避免了并发访问冲突。
val channel1 = Channel()
val channel2 = Channel()
val channel3 = Channel()
select {
channel1.onSend("message") { sendChannel: SendChannel ->
println("The message was successfully sent")
}
channel2.onReceive { message: String ->
println("The message was successfully received")
}
channel3.onReceiveCatching { result: ChannelResult ->
println("The message was successfully received")
}
}
import kotlinx.coroutines.selects.onTimeout
select {
onTimeout(5.seconds) {
println("Timeout")
}
}
fun main(): Unit = runBlocking {
var count = 0
val scope = CoroutineScope(Dispatchers.Default) // 使用多线程调度器
val job1 = scope.launch {
repeat(100_000_000) {
count++
}
}
val job2 = scope.launch {
repeat(100_000_000) {
count--
}
}
job1.join()
job2.join()
println("the final count is $count")
}
val myThreadLocal = ThreadLocal()
fun main(): Unit = runBlocking {
val scope = CoroutineScope(Dispatchers.Default)
val job = scope.launch {
myThreadLocal.set("hello")
println("Before: ${myThreadLocal.get()} on ${Thread.currentThread().name}")
delay(300)
// 可能切换到了另一个线程上
println("After: ${myThreadLocal.get()} on ${Thread.currentThread().name}")
}
// 尽可能抢占线程
repeat(30) {
scope.launch {
delay(500)
}
}
job.join()
}
Before: hello on DefaultDispatcher-worker-1
After: null on DefaultDispatcher-worker-19
val myThreadLocal = ThreadLocal()
fun main(): Unit = runBlocking {
val scope = CoroutineScope(Dispatchers.Default)
// 使用 asContextElement 将 ThreadLocal 包装成 CoroutineContext.Element
val job = scope.launch(myThreadLocal.asContextElement(value = "hello")) {
println("Before: ${myThreadLocal.get()} on ${Thread.currentThread().name}")
delay(1000)
// 无论切换到哪个线程,ThreadLocal 都会保持正确的值
println("After: ${myThreadLocal.get()} on ${Thread.currentThread().name}")
}
// 尽可能抢占线程
repeat(30) {
scope.launch {
delay(500)
}
}
job.join()
}
Before: hello on DefaultDispatcher-worker-1
After: hello on DefaultDispatcher-worker-17
Job
Job
select
Deferred
onAwait
async
Channel
onSend
onReceive
onReceiveCatching
Channel
Channel
Channel
select
onTimeout
ThreadLocal









