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

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

BufferOverflow 策略詳解
當(dāng) Channel 的緩沖區(qū)滿時(shí),BufferOverflow 參數(shù)決定了如何處理新的發(fā)送請(qǐng)求:
SUSPEND(默認(rèn)策略)
- 行為:當(dāng)緩沖區(qū)滿時(shí),掛起發(fā)送操作直到有空間可用
- 特點(diǎn):提供背壓控制,防止生產(chǎn)者過(guò)快
- 使用場(chǎng)景:需要確保所有數(shù)據(jù)都被處理的場(chǎng)景
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
- 行為:當(dāng)緩沖區(qū)滿時(shí),丟棄新元素,保留緩沖區(qū)中的現(xiàn)有元素
- 特點(diǎn):保護(hù)已緩沖的數(shù)據(jù)不被覆蓋
- 使用場(chǎng)景:保護(hù)重要的歷史數(shù)據(jù),防止新數(shù)據(jù)覆蓋
- 性能特點(diǎn):發(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)
}可以看到,當(dāng)緩沖區(qū)滿時(shí),會(huì)把新數(shù)據(jù)丟棄掉,因此,接收端只接收到了舊數(shù)據(jù)。

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

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

