Kotlin 協(xié)程之Channel的概念和基本使用詳解
前言
在 專欄 之前的文章中,我們已經(jīng)知道了協(xié)程的啟動、掛起、取消、異常以及常用的協(xié)程作用域等基礎應用。
這些基礎應用適合的場景是一次性任務,執(zhí)行完就結(jié)束了的場景。
launch / async 適合的場景
- 網(wǎng)絡請求
- 數(shù)據(jù)庫查詢
- 文件讀寫
- 并行計算任務
- 等等
而對于一些相對復雜的場景,例如:持續(xù)的數(shù)據(jù)流、需要在不同的協(xié)程之間傳遞數(shù)據(jù)、需要順序或背壓控制等場景,基礎的 launch / async
就不夠用了。
例如:
- 用戶點擊、輸入等事件流的處理
- 生產(chǎn)者-消費者模型的需求:任務排隊、日志流
- 高頻數(shù)據(jù)源處理(相機幀、音頻流等)
類似這種持續(xù)的、需要順序控制、或者多個協(xié)程配合執(zhí)行的場景,就需要用到 Channel 了。
Channel 的概念和基本使用
概念
顧名思義,Channel 有管道、通道的意思。Channel 跟 Java 中的 BlockingQueue 很相似,區(qū)別在于 Channel 是掛起的,不是阻塞的。
Channel 的核心特點就是能夠在不同的協(xié)程之間進行數(shù)據(jù)傳遞,并且能夠控制數(shù)據(jù)傳遞的順序。
使用起來很簡單,基本就分為以下幾步:
- 創(chuàng)建 Channel
- 通過
channel.send發(fā)送數(shù)據(jù) - 通過
channel.receive接收數(shù)據(jù)
整體的概念也比較簡單形象,就是一根管道,一個口子發(fā)送數(shù)據(jù),一個口子接收數(shù)據(jù)。
Channel 的創(chuàng)建
先來看下 Channel 的源碼,可以看到會根據(jù)傳入的參數(shù)選擇不同的實現(xiàn)。
public fun <E> Channel(
capacity: Int = RENDEZVOUS,
onBufferOverflow: BufferOverflow = BufferOverflow.SUSPEND,
onUndeliveredElement: ((E) -> Unit)? = null
): Channel<E> =
when (capacity) {
RENDEZVOUS -> {
if (onBufferOverflow == BufferOverflow.SUSPEND)
BufferedChannel(RENDEZVOUS, onUndeliveredElement) // an efficient implementation of rendezvous channel
else
ConflatedBufferedChannel(
1,
onBufferOverflow,
onUndeliveredElement
) // support buffer overflow with buffered channel
}
CONFLATED -> {
require(onBufferOverflow == BufferOverflow.SUSPEND) {
"CONFLATED capacity cannot be used with non-default onBufferOverflow"
}
ConflatedBufferedChannel(1, BufferOverflow.DROP_OLDEST, onUndeliveredElement)
}
UNLIMITED -> BufferedChannel(
UNLIMITED,
onUndeliveredElement
) // ignores onBufferOverflow: it has buffer, but it never overflows
BUFFERED -> { // uses default capacity with SUSPEND
if (onBufferOverflow == BufferOverflow.SUSPEND) BufferedChannel(
CHANNEL_DEFAULT_CAPACITY,
onUndeliveredElement
)
else ConflatedBufferedChannel(1, onBufferOverflow, onUndeliveredElement)
}
else -> {
if (onBufferOverflow === BufferOverflow.SUSPEND) BufferedChannel(capacity, onUndeliveredElement)
else ConflatedBufferedChannel(capacity, onBufferOverflow, onUndeliveredElement)
}
}參數(shù)概覽
| 參數(shù) | 類型 | 默認值 | 描述 |
|---|---|---|---|
capacity | Int | RENDEZVOUS | 通道容量,決定緩沖區(qū)大小和行為模式 |
onBufferOverflow | BufferOverflow | SUSPEND | 緩沖區(qū)溢出時的處理策略 |
onUndeliveredElement | ((E) -> Unit)? | null | 元素未能送達時的回調(diào)函數(shù) |
capacity(容量配置)
capacity 參數(shù)決定了 Channel 的緩沖行為和容量大?。?/p>
RENDEZVOUS(值為 0):無緩沖,發(fā)送者和接收者必須同時準備好CONFLATED(值為 -1):只保留最新的元素,舊元素會被覆蓋UNLIMITED(值為Int.MAX_VALUE):理論上就是無限容量,永不阻塞發(fā)送BUFFERED(值為 64):默認緩沖大小- 自定義正整數(shù):自己指定具體的緩沖區(qū)大小
onBufferOverflow(溢出策略)
當緩沖區(qū)滿時的處理策略:
SUSPEND:掛起發(fā)送操作,等待緩沖區(qū)有空間(默認)DROP_OLDEST:丟棄舊的元素,添加新元素DROP_LATEST:丟棄新元素,保留緩沖區(qū)中的現(xiàn)有元素
onUndeliveredElement(未送達回調(diào))
當元素無法送達時的清理回調(diào)函數(shù):
null:不執(zhí)行任何清理操作(默認)- 自定義函數(shù):用于資源清理、日志記錄等,根據(jù)業(yè)務需求來定義
參數(shù)組合效果
| capacity | onBufferOverflow | 行為 | 適用場景 |
|---|---|---|---|
RENDEZVOUS | SUSPEND | 無緩沖,同步通信 | 嚴格的生產(chǎn)者-消費者同步 |
BUFFERED | SUSPEND | 有限緩沖,滿時掛起 | 一般的異步處理,默認的緩沖數(shù)量是 64 |
UNLIMITED | SUSPEND | 緩沖長度為 Int.MAX_VALUE | 高吞吐量場景(生產(chǎn)上不建議使用,有內(nèi)存方面的風險) |
CONFLATED | DROP_OLDEST | 無緩沖,只保留最新值 | 狀態(tài)更新、實時數(shù)據(jù) |
| 自定義大小 | SUSPEND | 固定大小,滿時掛起 | 批量處理、批量任務 |
| 自定義大小 | DROP_OLDEST | 固定大小,丟棄舊數(shù)據(jù) | 獲取最近 N 個元素 |
| 自定義大小 | DROP_LATEST | 固定大小,拒絕新數(shù)據(jù) | 保護重要歷史數(shù)據(jù) |
Capacity
RENDEZVOUS(會合模式)
特點:
- 容量為 0,無緩沖區(qū)
- 發(fā)送者和接收者必須同時準備好才能完成數(shù)據(jù)傳輸
- 提供強同步保證,一手交錢一手交貨
使用示例:
suspend fun demonstrateRendezvousChannel() {
// 創(chuàng)建 RENDEZVOUS Channel(默認容量為 0),默認什么都不傳就是 rendezvous 模式,Channel<String>()
val rendezvousChannel = Channel<String>(Channel.RENDEZVOUS)
// 啟動發(fā)送者協(xié)程
val senderJob = GlobalScope.launch {
println("[發(fā)送者] 準備發(fā)送消息...")
rendezvousChannel.send("Hello from RENDEZVOUS!")
println("[發(fā)送者] 消息已發(fā)送")
rendezvousChannel.send("Second message")
println("[發(fā)送者] 第二條消息已發(fā)送")
rendezvousChannel.close()
}
// 啟動接收者協(xié)程
val receiverJob = GlobalScope.launch {
delay(1000) // 延遲1秒,發(fā)送者會等待接收者準備好
println("[接收者] 開始接收消息...")
for (message in rendezvousChannel) {
println("[接收者] 收到消息: $message")
delay(500) // 模擬處理時間
}
println("[接收者] Channel已關閉")
}
// 等待所有協(xié)程完成
joinAll(senderJob, receiverJob)
}執(zhí)行結(jié)果

CONFLATED(只留最新值)
特點:
- 容量為 1,但會丟棄舊值
- 只保留最新的元素
- 發(fā)送操作永不阻塞
- 只能使用
BufferOverflow.SUSPEND策略
源碼分析:
CONFLATED -> {
require(onBufferOverflow == BufferOverflow.SUSPEND) {
"CONFLATED capacity cannot be used with non-default onBufferOverflow"
}
ConflatedBufferedChannel(1, BufferOverflow.DROP_OLDEST, onUndeliveredElement)
}使用示例:
suspend fun demonstrateConflatedChannel() {
// 創(chuàng)建 CONFLATED Channel,相當于:Channel<String>(1, BufferOverflow.DROP_OLDEST)
val conflatedChannel = Channel<String>(Channel.CONFLATED)
// 快速發(fā)送多個消息
val senderJob = GlobalScope.launch {
repeat(5) { i ->
val message = "Update-$i"
conflatedChannel.send(message)
println("[發(fā)送者] 發(fā)送更新: $message")
delay(100) // 短暫延遲
}
conflatedChannel.close()
}
// 慢速接收者
val receiverJob = GlobalScope.launch {
delay(1000) // 延遲1秒,讓發(fā)送者發(fā)送完所有消息
println("[接收者] 開始接收(只會收到最新的值)...")
for (message in conflatedChannel) {
println("[接收者] 收到: $message")
}
}
joinAll(senderJob, receiverJob)
}
UNLIMITED(無限容量)
特點:
- 容量為
Int.MAX_VALUE,理論上無限容量 - 發(fā)送操作永不阻塞,但要注意內(nèi)存使用
- 忽略
onBufferOverflow參數(shù) - 適用于高吞吐量場景,但生產(chǎn)環(huán)境需謹慎使用
suspend fun demonstrateUnlimitedChannel() {
val unlimitedChannel = Channel<String>(Channel.UNLIMITED)
val senderJob = GlobalScope.launch {
repeat(10) { i ->
val message = "Message-$i"
unlimitedChannel.send(message)
println("[發(fā)送者] 立即發(fā)送: $message")
}
unlimitedChannel.close()
println("[發(fā)送者] 所有消息已發(fā)送,Channel已關閉")
}
val receiverJob = GlobalScope.launch {
delay(1000) // 延遲1秒開始接收
println("[接收者] 開始慢速接收...")
for (message in unlimitedChannel) {
println("[接收者] 處理: $message")
delay(300) // 模擬處理時間
}
}
joinAll(senderJob, receiverJob)
}
BUFFERED(有限容量)
特點:
- 使用默認容量 (
CHANNEL_DEFAULT_CAPACITY,通常為 64) - 在緩沖區(qū)滿時根據(jù)
onBufferOverflow策略處理
源碼分析:
BUFFERED -> {
if (onBufferOverflow == BufferOverflow.SUSPEND)
BufferedChannel(CHANNEL_DEFAULT_CAPACITY, onUndeliveredElement)
else
ConflatedBufferedChannel(1, onBufferOverflow, onUndeliveredElement)
}使用示例:
suspend fun demonstrateBufferedDefaultChannel() {
// 創(chuàng)建 BUFFERED Channel(默認容量為 64)
val bufferedChannel = Channel<String>(Channel.BUFFERED)
val senderJob = GlobalScope.launch {
repeat(100) { i ->
bufferedChannel.send("Message-$i")
println("[發(fā)送者] 發(fā)送 Message-$i")
}
bufferedChannel.close()
}
val receiverJob = GlobalScope.launch {
delay(1000) // 延遲接收
for (message in bufferedChannel) {
println("[接收者] 收到: $message")
delay(50)
}
}
joinAll(senderJob, receiverJob)
}與下面自定義容量效果類似。
自定義容量
特點:
- 指定具體的緩沖區(qū)大小
- 根據(jù)
onBufferOverflow策略處理溢出
源碼分析:
else -> {
if (onBufferOverflow === BufferOverflow.SUSPEND)
BufferedChannel(capacity, onUndeliveredElement)
else
ConflatedBufferedChannel(capacity, onBufferOverflow, onUndeliveredElement)
}使用示例:
suspend fun demonstrateBufferedChannel() {
// 創(chuàng)建容量為3的緩沖Channel
val bufferedChannel = Channel<Int>(capacity = 3)
// 啟動發(fā)送者協(xié)程
val senderJob = GlobalScope.launch {
repeat(5) { i ->
println("[發(fā)送者] 發(fā)送數(shù)字: $i")
bufferedChannel.send(i)
println("[發(fā)送者] 數(shù)字 $i 已發(fā)送")
}
bufferedChannel.close()
println("[發(fā)送者] Channel已關閉")
}
// 啟動接收者協(xié)程,延遲接收以觀察緩沖效果
val receiverJob = GlobalScope.launch {
delay(2000) // 延遲2秒開始接收
println("[接收者] 開始接收數(shù)字...")
for (number in bufferedChannel) {
println("[接收者] 收到數(shù)字: $number")
delay(800) // 模擬慢速處理
}
}
joinAll(senderJob, receiverJob)
}可以看到,因為默認的溢出策略是 SUSPEND,所以當緩沖區(qū)滿了時,發(fā)送者會被掛起,直到接收者處理完一個元素,才會繼續(xù)發(fā)送。

BufferOverflow 策略詳解
當 Channel 的緩沖區(qū)滿時,BufferOverflow 參數(shù)決定了如何處理新的發(fā)送請求:
SUSPEND(默認策略)
- 行為:當緩沖區(qū)滿時,掛起發(fā)送操作直到有空間可用
- 特點:提供背壓控制,防止生產(chǎn)者過快
- 使用場景:需要確保所有數(shù)據(jù)都被處理的場景
suspend fun demonstrateBasicOperations() {
//容量為 2,溢出策略為SUSPEND
val channel = Channel<String>(capacity = 2, onBufferOverflow = BufferOverflow.SUSPEND)
//發(fā)送的速度快
val job1 = GlobalScope.launch {
repeat(5) {
channel.send("Message-$it")
println("[發(fā)送者] 發(fā)送 Message-$it")
}
channel.close()
}
val job2 = GlobalScope.launch {
//除了用 channel.recrive 外,也可以直接 用 for 循環(huán)接收數(shù)據(jù)
for (message in channel) {
//接收的速度慢
delay(1000)
println("[接收者] 接收到: $message")
}
}
joinAll(job1, job2)
}
DROP_LATEST
- 行為:當緩沖區(qū)滿時,丟棄新元素,保留緩沖區(qū)中的現(xiàn)有元素
- 特點:保護已緩沖的數(shù)據(jù)不被覆蓋
- 使用場景:保護重要的歷史數(shù)據(jù),防止新數(shù)據(jù)覆蓋
- 性能特點:發(fā)送操作永不阻塞,但新數(shù)據(jù)可能被丟棄
suspend fun demonstrateBasicOperations() {
val channel = Channel<String>(capacity = 2, onBufferOverflow = BufferOverflow.DROP_LATEST)
val job1 = GlobalScope.launch {
repeat(5) {
channel.send("Message-$it")
println("[發(fā)送者] 發(fā)送 Message-$it")
}
channel.close()
}
val job2 = GlobalScope.launch {
for (message in channel) {
delay(1000)
println("[接收者] 接收到: $message")
}
}
joinAll(job1, job2)
}可以看到,當緩沖區(qū)滿時,會把新數(shù)據(jù)丟棄掉,因此,接收端只接收到了舊數(shù)據(jù)。

DROP_OLDEST
- 行為:當緩沖區(qū)滿時,丟棄舊的元素,添加新元素
- 特點:保持固定的內(nèi)存使用,優(yōu)先保留新數(shù)據(jù)
- 使用場景:實時數(shù)據(jù)流、最近N個元素
- 性能特點:發(fā)送操作永不阻塞,但可能丟失歷史數(shù)據(jù)
suspend fun demonstrateBasicOperations() {
val channel = Channel<String>(capacity = 2, onBufferOverflow = BufferOverflow.DROP_OLDEST)
val job1 = GlobalScope.launch {
repeat(5) {
channel.send("Message-$it")
println("[發(fā)送者] 發(fā)送 Message-$it")
}
channel.close()
}
val job2 = GlobalScope.launch {
for (message in channel) {
delay(1000)
println("[接收者] 接收到: $message")
}
}
joinAll(job1, job2)
}
需要注意的是,當緩沖區(qū)滿了之后,1 和 2 被丟棄了,3 和 4 被放進去了。從這里可以看出,丟棄數(shù)據(jù)時,并不是把最早的舊數(shù)據(jù)丟掉,這里跟內(nèi)部的實現(xiàn)有關。
onUndeliveredElement 回調(diào)
當元素無法送達時(如 Channel 被取消或關閉),會調(diào)用此回調(diào)函數(shù)
suspend fun demonstrateBasicOperations() {
val channel = Channel<String>(capacity = 2, onBufferOverflow = BufferOverflow.DROP_OLDEST) {
println("[Channel] 緩沖區(qū)已滿,無法放到緩沖區(qū),值:${it}")
}
// 演示基本的send和receive操作
val job1 = GlobalScope.launch {
repeat(5) {
channel.send("Message-$it")
println("[發(fā)送者] 發(fā)送 Message-$it")
}
channel.close()
}
val job2 = GlobalScope.launch {
for (message in channel) {
delay(1000)
println("[接收者] 接收到: $message")
}
}
joinAll(job1, job2)
}
Channel 操作方式
Channel 提供了兩種操作方式:阻塞操作和非阻塞操作。
阻塞操作(send/receive)
send() 和 receive() 方法都是掛起方法,它們會阻塞當前協(xié)程,直到完成操作。
非阻塞操作(trySend/tryReceive)
trySend() 和 tryReceive() 是 Channel 提供的非阻塞操作 API。與阻塞版本不同,這些方法會立即返回結(jié)果,不會掛起當前協(xié)程,也不會拋出異常。
操作對比
| 操作類型 | 阻塞版本 | 非阻塞版本 | 行為差異 |
|---|---|---|---|
| 發(fā)送 | send() | trySend() | send() 會掛起直到有空間;trySend() 立即返回結(jié)果 |
| 接收 | receive() | tryReceive() | receive() 會掛起直到有數(shù)據(jù);tryReceive() 立即返回結(jié)果 |
返回值類型
trySend()返回ChannelResult<Unit>tryReceive()返回ChannelResult<T>
ChannelResult 是一個密封類,通過密封類中的成員 isSuccess 和 getOrNull() 可以判斷操作是否成功。

大部分場景下,send / receive + 合理的 Channel 配置就能解決問題,trySend/tryReceive 更多的是想達到如下效果:
- 避免不必要的協(xié)程掛起開銷,希望立即得到結(jié)果
- 提供更精細的控制邏輯,如:超時處理、重試機制等
- 實現(xiàn)更好的錯誤處理和用戶反饋,能更好地處理異常場景
runBlocking {
val channel = Channel<Int>(2)
val sendJob = launch {
repeat(5) {
delay(100)
val sendResult = channel.trySend(it)
sendResult.onSuccess {
println("發(fā)送成功")
}.onFailure {
println("發(fā)送失敗")
}.onClosed {
println("通道已關閉")
}
}
}
val receiveJob = launch {
for (i in channel) {
delay(300)
println("接收到數(shù)據(jù):${i}")
}
}
joinAll(sendJob, receiveJob)
}
Channel 狀態(tài)管理
Channel 在其生命周期中會經(jīng)歷以下幾個關鍵狀態(tài):
- 活躍狀態(tài)(Active):可以正常發(fā)送和接收數(shù)據(jù)
- 發(fā)送端關閉(Closed for Send):不能發(fā)送新數(shù)據(jù),但可以接收緩沖區(qū)中的數(shù)據(jù)
- 接收端關閉(Closed for Receive):不能接收數(shù)據(jù),緩沖區(qū)已清空
- 取消狀態(tài)(Cancelled):Channel 被取消,所有操作都會失敗
API
channel.close():關閉 Channelchannel.isClosedForSend:判斷發(fā)送端是否已關閉channel.isClosedForReceive:判斷接收端是否已關閉channel.cancel():取消 Channel
Close(關閉操作)
- 調(diào)用
close()后,isClosedForSend立即變?yōu)?true - 此時,緩沖區(qū)中的數(shù)據(jù)仍可被消費
- 只有當緩沖區(qū)清空后,
isClosedForReceive才變?yōu)?true
示例:
suspend fun demonstrateChannelClose() {
val channel = Channel<String>(1)
val producer = GlobalScope.launch {
try {
for (i in 1..5) {
val message = "Message $i"
println("準備發(fā)送: $message")
channel.send(message)
println("成功發(fā)送: $message")
delay(100)
}
} catch (e: ClosedSendChannelException) {
println("生產(chǎn)者: Channel已關閉,無法發(fā)送數(shù)據(jù) - ${e.message}")
}
}
val consumer = GlobalScope.launch {
try {
for (message in channel) {
println("接收到: $message")
delay(200)
}
println("消費者: Channel已關閉,退出接收循環(huán)")
} catch (e: Exception) {
println("消費者異常: ${e.message}")
}
}
delay(300) // 模擬讓一些數(shù)據(jù)能夠被接收到
// 檢查Channel狀態(tài)
println("關閉前狀態(tài):")
println(" isClosedForSend: ${channel.isClosedForSend}")
println(" isClosedForReceive: ${channel.isClosedForReceive}")
// 關閉Channel
println("\n正在關閉Channel...")
channel.close()
// 檢查關閉后的狀態(tài)
println("關閉后狀態(tài):")
println(" isClosedForSend: ${channel.isClosedForSend}")
println(" isClosedForReceive: ${channel.isClosedForReceive}")
// 等待協(xié)程完成
producer.join()
consumer.join()
println("最終狀態(tài):")
println(" isClosedForSend: ${channel.isClosedForSend}")
println(" isClosedForReceive: ${channel.isClosedForReceive}")
}
Cancel(取消操作)
cancel() 方法用于強制取消 Channel,它會:
- 立即關閉發(fā)送和接收端
- 清空緩沖區(qū)中的所有數(shù)據(jù)
- 觸發(fā)
onUndeliveredElement回調(diào)(如果設置了)
suspend fun demonstrateChannelCancel() {
val channel = Channel<String>(capacity = 5) {
println("消息未被接收:${it}")
}
val producer = GlobalScope.launch {
try {
for (i in 1..8) {
val message = "Message $i"
println("嘗試發(fā)送: $message")
channel.send(message)
println("成功發(fā)送: $message")
delay(100)
}
} catch (e: CancellationException) {
println("生產(chǎn)者: Channel被取消 - ${e.message}")
}
}
val consumer = GlobalScope.launch {
try {
for (message in channel) {
println("接收到: $message")
delay(300)
}
} catch (e: CancellationException) {
println("消費者: 協(xié)程被取消 - ${e.message}")
}
}
delay(400) // 讓一些操作執(zhí)行
println("\n取消前狀態(tài):")
println(" isClosedForSend: ${channel.isClosedForSend}")
println(" isClosedForReceive: ${channel.isClosedForReceive}")
// 取消Channel
println("\n正在取消Channel...")
channel.cancel(CancellationException("主動取消Channel"))
println("取消后狀態(tài):")
println(" isClosedForSend: ${channel.isClosedForSend}")
println(" isClosedForReceive: ${channel.isClosedForReceive}")
// 等待協(xié)程完成
producer.join()
consumer.join()
}
Channel 異常處理
在使用 Channel 的過程中,會遇到各種異常情況。主要包括以下幾種類型:
ClosedSendChannelException
觸發(fā)條件:
- 在已關閉的 Channel 上調(diào)用
send()方法 - Channel 調(diào)用
close()后,發(fā)送端立即關閉
示例:
suspend fun demonstrateClosedSendException() {
val channel = Channel<String>()
// 關閉 Channel
channel.close()
try {
// 嘗試在已關閉的 Channel 上發(fā)送數(shù)據(jù)
channel.send("This will throw exception")
} catch (e: ClosedSendChannelException) {
println("捕獲異常: ${e.message}")
println("異常類型: ${e::class.simpleName}")
}
}ClosedReceiveChannelException
觸發(fā)條件:
- 從已關閉且緩沖區(qū)為空的 Channel 調(diào)用
receive()方法 - 當
isClosedForReceive為true時調(diào)用receive()
示例:
suspend fun demonstrateClosedReceiveException() {
val channel = Channel<String>()
// 關閉 Channel
channel.close()
try {
// 嘗試從已關閉且空的 Channel 接收數(shù)據(jù)
val message = channel.receive()
println("收到消息: $message")
} catch (e: ClosedReceiveChannelException) {
println("捕獲異常: ${e.message}")
println("異常類型: ${e::class.simpleName}")
}
}CancellationException
觸發(fā)條件:
- Channel 被
cancel()方法取消 - 父協(xié)程被取消,導致 Channel 操作被取消
- 超時或其他取消信號
示例:
suspend fun demonstrateCancellationException() {
val channel = Channel<String>()
val job = GlobalScope.launch {
try {
// 這個操作會被取消
channel.send("This will be cancelled")
} catch (e: CancellationException) {
println("發(fā)送操作被取消: ${e.message}")
throw e // 重新拋出 CancellationException
}
}
delay(100)
// 取消 Channel
channel.cancel(CancellationException("手動取消 Channel"))
try {
job.join()
} catch (e: CancellationException) {
println("協(xié)程被取消: ${e.message}")
}
}異常與狀態(tài)關系
| Channel 狀態(tài) | send() 行為 | receive() 行為 | trySend() 行為 | tryReceive() 行為 |
|---|---|---|---|---|
| 活躍狀態(tài) | 正常發(fā)送或掛起 | 正常接收或掛起 | 返回成功/失敗結(jié)果 | 返回成功/失敗結(jié)果 |
| 發(fā)送端關閉 | 拋出 ClosedSendChannelException | 正常接收緩沖區(qū)數(shù)據(jù) | 返回失敗結(jié)果 | 正常返回結(jié)果 |
| 接收端關閉 | 拋出 ClosedSendChannelException | 拋出 ClosedReceiveChannelException | 返回失敗結(jié)果 | 返回失敗結(jié)果 |
| 已取消 | 拋出 CancellationException | 拋出 CancellationException | 返回失敗結(jié)果 | 返回失敗結(jié)果 |
異常處理技巧
使用非阻塞操作避免異常
非阻塞操作不會拋出異常,而是返回結(jié)果對象:
suspend fun safeChannelOperations() {
val channel = Channel<String>()
// 安全的發(fā)送操作
val sendResult = channel.trySend("Safe message")
when {
sendResult.isSuccess -> println("發(fā)送成功")
sendResult.isFailure -> println("發(fā)送失敗: ${sendResult.exceptionOrNull()}")
sendResult.isClosed -> println("Channel 已關閉")
}
// 安全的接收操作
val receiveResult = channel.tryReceive()
when {
receiveResult.isSuccess -> println("接收到: ${receiveResult.getOrNull()}")
receiveResult.isFailure -> println("接收失敗: ${receiveResult.exceptionOrNull()}")
receiveResult.isClosed -> println("Channel 已關閉")
}
}健壯的異常處理
suspend fun robustChannelUsage() {
val channel = Channel<String>()
val producer = GlobalScope.launch {
try {
repeat(5) { i ->
if (channel.isClosedForSend) {
println("Channel 已關閉,停止發(fā)送")
break
}
channel.send("Message $i")
delay(100)
}
} catch (e: ClosedSendChannelException) {
println("生產(chǎn)者: Channel 已關閉")
} catch (e: CancellationException) {
println("生產(chǎn)者: 操作被取消")
throw e // 重新拋出取消異常
} finally {
println("生產(chǎn)者: 清理資源")
}
}
val consumer = GlobalScope.launch {
try {
while (!channel.isClosedForReceive) {
try {
val message = channel.receive()
println("消費者: 收到 $message")
} catch (e: ClosedReceiveChannelException) {
println("消費者: Channel 已關閉且無更多數(shù)據(jù)")
break
}
delay(200)
}
} catch (e: CancellationException) {
println("消費者: 操作被取消")
throw e
} finally {
println("消費者: 清理資源")
}
}
delay(1000)
channel.close()
joinAll(producer, consumer)
}總結(jié)
Channel 關鍵概念對比
| 特性 | RENDEZVOUS | CONFLATED | BUFFERED | UNLIMITED | 自定義容量 |
|---|---|---|---|---|---|
| 容量 | 0 | 1 | 64 | Int.MAX_VALUE | 指定值 |
| 緩沖行為 | 無緩沖,同步 | 只保留最新值 | 有限緩沖 | 無限緩沖 | 有限緩沖 |
| 發(fā)送阻塞 | 是 | 否 | 緩沖滿時 | 否 | 緩沖滿時 |
| 適用場景 | 嚴格同步 | 狀態(tài)更新 | 一般異步 | 高吞吐量 | 批量處理 |
| 內(nèi)存風險 | 低 | 低 | 中等 | 高 | 可控 |
溢出策略對比
| 策略 | 行為 | 性能特點 | 適用場景 |
|---|---|---|---|
| SUSPEND | 掛起發(fā)送操作 | 提供背壓控制 | 確保數(shù)據(jù)完整性 |
| DROP_OLDEST | 丟棄舊元素 | 發(fā)送不阻塞 | 實時數(shù)據(jù)流 |
| DROP_LATEST | 丟棄新元素 | 發(fā)送不阻塞 | 保護歷史數(shù)據(jù) |
操作方式
| 操作類型 | 阻塞版本 | 非阻塞版本 | 異常處理 | 返回值 |
|---|---|---|---|---|
| 發(fā)送 | send() | trySend() | 拋出異常 | ChannelResult<Unit> |
| 接收 | receive() | tryReceive() | 拋出異常 | ChannelResult<T> |
| 特點 | 會掛起協(xié)程 | 立即返回 | 需要 try-catch | 通過結(jié)果對象判斷 |
Channel 狀態(tài)生命周期
| 狀態(tài) | 描述 | send() | receive() | 檢查方法 |
|---|---|---|---|---|
| 活躍 | 正常工作狀態(tài) | ? 正常 | ? 正常 | - |
| 發(fā)送關閉 | 調(diào)用 close() 后 | ? 異常 | ? 可接收緩沖區(qū)數(shù)據(jù) | isClosedForSend |
| 接收關閉 | 緩沖區(qū)清空后 | ? 異常 | ? 異常 | isClosedForReceive |
| 已取消 | 調(diào)用 cancel() 后 | ? 異常 | ? 異常 | - |
總體來說,Channel 是一種非常強大的協(xié)程通信機制,它可以幫助我們在協(xié)程之間進行安全、高效的通信。在使用 Channel時,我們需要注意異常處理、緩沖區(qū)容量、溢出策略等問題。
感謝閱讀,如果對你有幫助請三連(點贊、收藏、加關注)支持。有任何疑問或建議,歡迎在評論區(qū)留言討論。如需轉(zhuǎn)載,請注明出處:喻志強的博客
到此這篇關于Kotlin 協(xié)程之Channel的概念和基本使用詳解的文章就介紹到這了,更多相關Kotlin 協(xié)程Channel使用內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
Android之高德地圖定位SDK集成及地圖功能實現(xiàn)
本文主要介紹了Android中高德地圖定位SDK集成及地圖功能的實現(xiàn)。具有很好的參考價值。下面跟著小編一起來看下吧2017-04-04
RxJava 1升級到RxJava 2過程中踩過的一些“坑”
RxJava2相比RxJava1,它的改動還是很大的,那么下面這篇文章主要給大家總結(jié)了在RxJava 1升級到RxJava 2過程中踩過的一些“坑”,文中介紹的非常詳細,對大家具有一定的參考學習價值,需要的朋友們下來要一起看看吧。2017-05-05
Android開發(fā)仿QQ空間根據(jù)位置彈出PopupWindow顯示更多操作效果
我們打開QQ空間的時候有個箭頭按鈕點擊之后彈出PopupWindow會根據(jù)位置的變化顯示在箭頭的上方還是下方,比普通的PopupWindow彈在屏幕中間顯示好看的多,今天就給大家分享下實現(xiàn)代碼,需要的朋友參考下吧2016-12-12
Android Studio3.6中的View Binding初探及用法區(qū)別
這篇文章主要介紹了Android 中的View Binding初探及用法區(qū)別,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下2020-03-03

