You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让Kotlin协程Channel按指定顺序接收,保证任务Step3按序执行?

解决方案:实现按初始顺序执行第三步

核心问题在于第二步并发完成后,任务会无序进入第三步的处理流程,导致第三步执行顺序混乱。我们需要引入有序缓冲区,确保第三步严格按照任务的初始索引顺序执行,不受第二步完成先后的影响。

修改后的代码实现

import android.util.Log
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlin.random.Random
import java.util.concurrent.ConcurrentSkipListMap

class SortingReceiveChannelTest(
    private val coroutineScope: CoroutineScope
) {
    private val tasks = List(30) { Task(it) }

    private val randomDelay = Random(Int.MAX_VALUE)
    private val step2Channel = Channel<Task>(2)
    // 用有序映射存储完成step2的任务,key为任务索引,自动维护顺序
    private val completedStep2Tasks = ConcurrentSkipListMap<Int, Task>()
    private val mutex = Mutex()
    private var currentStep3Index = 0

    fun start() {
        coroutineScope.launch(Dispatchers.IO) {
            for (task in tasks) {
                task.step1()
                step2Channel.send(task)
            }
        }

        coroutineScope.launch(Dispatchers.IO) {
            repeat(tasks.size) {
                val task = step2Channel.receive()
                coroutineScope.launch(Dispatchers.IO) {
                    task.step2()
                    // step2完成后存入有序映射,触发第三步顺序检查
                    completedStep2Tasks[task.index] = task
                    tryProcessStep3()
                }
            }
        }
    }

    private suspend fun tryProcessStep3() {
        mutex.withLock {
            // 循环检查当前待执行的索引是否已完成step2,直到没有连续的已完成任务
            while (completedStep2Tasks.containsKey(currentStep3Index)) {
                val task = completedStep2Tasks.remove(currentStep3Index)!!
                task.step3()
                currentStep3Index++
            }
        }
    }

    inner class Task(val index: Int) {
        fun step1() {
            Log.i(TAG, "execute step1[$index] sequential")
            Thread.sleep(randomDelay.nextLong(100, 200))
        }

        fun step2() {
            Log.i(TAG, "execute step2[$index] concurrently")
            if (index == 2 || index == 5 || index == 8 || index == 16) {
                Thread.sleep(randomDelay.nextLong(2000, 5000))
            } else {
                Thread.sleep(randomDelay.nextLong(100, 300))
            }
        }

        fun step3() {
            Log.i(TAG, "execute step3[$index] sequential")
            Thread.sleep(randomDelay.nextLong(100, 300))
        }
    }

    companion object {
        private const val TAG = "SortChannel"
    }
}

关键逻辑说明

  • 有序存储:ConcurrentSkipListMap会自动按任务索引排序,让我们能快速定位当前需要执行的任务,同时支持并发读写。
  • 顺序触发:每个任务完成step2后,立即触发第三步的顺序检查逻辑,只要当前待执行的索引对应的任务已完成,就执行step3并递增索引,确保连续执行。
  • 线程安全:用Mutex保证索引修改和映射访问的原子性,避免并发场景下的逻辑冲突。

替代方案:基于Channel的排序中转

如果坚持使用Channel,可以在第三步前增加一个排序协程,将无序的任务重新排序后再发送到有序Channel:

// 新增有序Channel和排序协程
private val step3Channel = Channel<Task>()
private val step3OrderedChannel = Channel<Task>()

fun start() {
    // ... 原有step1、step2逻辑不变 ...

    // 排序中转协程
    coroutineScope.launch(Dispatchers.IO) {
        val sortedBuffer = mutableMapOf<Int, Task>()
        var currentIndex = 0
        repeat(tasks.size) {
            val task = step3Channel.receive()
            sortedBuffer[task.index] = task
            // 输出连续的已完成任务到有序Channel
            while (sortedBuffer.containsKey(currentIndex)) {
                step3OrderedChannel.send(sortedBuffer.remove(currentIndex)!!)
                currentIndex++
            }
        }
        step3OrderedChannel.close()
    }

    // 第三步从有序Channel接收任务
    coroutineScope.launch(Dispatchers.IO) {
        for (task in step3OrderedChannel) {
            task.step3()
        }
    }
}

这个思路和有序映射方案逻辑一致,只是通过多一层Channel实现任务的有序中转。

内容的提问来源于stack exchange,提问作者VinceStyling

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 13:13:08