最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Kotlin 協(xié)程之Channel的概念和基本使用詳解

 更新時間:2025年08月22日 09:34:17   作者:XeonYu  
文章介紹協(xié)程在復雜場景中使用Channel進行數(shù)據(jù)傳遞與控制,涵蓋創(chuàng)建參數(shù)、緩沖策略、操作方式及異常處理,適用于持續(xù)數(shù)據(jù)流、多協(xié)程協(xié)作等,需注意容量配置和狀態(tài)管理,本文給大家介紹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ù)傳遞的順序。
使用起來很簡單,基本就分為以下幾步:

  1. 創(chuàng)建 Channel
  2. 通過 channel.send 發(fā)送數(shù)據(jù)
  3. 通過 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ù)類型默認值描述
capacityIntRENDEZVOUS通道容量,決定緩沖區(qū)大小和行為模式
onBufferOverflowBufferOverflowSUSPEND緩沖區(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ù)組合效果

capacityonBufferOverflow行為適用場景
RENDEZVOUSSUSPEND無緩沖,同步通信嚴格的生產(chǎn)者-消費者同步
BUFFEREDSUSPEND有限緩沖,滿時掛起一般的異步處理,默認的緩沖數(shù)量是 64
UNLIMITEDSUSPEND緩沖長度為 Int.MAX_VALUE高吞吐量場景(生產(chǎn)上不建議使用,有內(nèi)存方面的風險)
CONFLATEDDROP_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 是一個密封類,通過密封類中的成員 isSuccessgetOrNull() 可以判斷操作是否成功。

大部分場景下,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():關閉 Channel
  • channel.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() 方法
  • isClosedForReceivetrue 時調(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 關鍵概念對比

特性RENDEZVOUSCONFLATEDBUFFEREDUNLIMITED自定義容量
容量0164Int.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ù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

最新評論

柯坪县| 江北区| 常山县| 高密市| 同德县| 衡水市| 麻阳| 临安市| 依安县| 宁晋县| 滦南县| 吉林省| 纳雍县| 荆门市| 安阳县| 厦门市| 洛川县| 舒城县| 富川| 拉萨市| 章丘市| 寻甸| 镇康县| 孙吴县| 达孜县| 黄骅市| 衡南县| 河北区| 通许县| 广河县| 朝阳市| 扶沟县| 阿拉善左旗| 宝清县| 虞城县| 大理市| 石林| 阜阳市| 台前县| 磐安县| 诸城市|