根据 官方文档可知select 表达式可以同时等待多个挂起函数,并选择第一个可用的。select具有一下功能可以细分如下:

  1. 在多个Channel中选择第一个到达的进行处理;
  2. 可以对主-从Channel模式进行分流处理,以降低主Channel的压力。

1. Channel的多路复用

下面通过官方的例子来分析Channel的多路复用原理,首先定义两个以不同速率产生字符串的生产者:

fun CoroutineScope.fizz() = produce<String> {
    while (true) { // sends "Fizz" every 300 ms
        delay(400)
        send("Fizz")
    }
}

fun CoroutineScope.buzz() = produce<String> {
    while (true) { // sends "Buzz!" every 500 ms
        delay(700)
        send("Buzz!")
    }
}

然后对以上两个Channel进行多路复用:

suspend fun selectFizzBuzz(fizz: ReceiveChannel<String>, buzz: ReceiveChannel<String>) {
    select<Unit> {
        fizz.onReceive { value ->  // this is the first select clause
            println("fizz -> '$value'")
        }
        buzz.onReceive { value ->  // this is the second select clause
            println("buzz -> '$value'")
        }
    }
}

下面是主函数在runBlocking函数中来开启所有操作:

private fun selectOne() {
    runBlocking<Unit> {
        val fizz = fizz()
        val buzz = buzz()
        repeat(7) {
            selectFizzBuzz(fizz, buzz)
        }
        coroutineContext.cancelChildren() // cancel fizz & buzz coroutines
    }
}

对于函数打印的结果如图所示,经过多路复用后两个Channel的数据以时间顺序被依次消费掉。
在这里插入图片描述

1.1. onReceive()函数细节

在select()函数中,可以先在参数SelectBuilder范围内使用_clauses_指定多个挂起函数,然后同时等待结果,直到某一个clauses变为_selected_或者_fails_状态为止。其源码如下:

public suspend inline fun <R> select(crossinline builder: SelectBuilder<R>.() -> Unit): R =
    suspendCoroutineUninterceptedOrReturn { uCont ->
        val scope = SelectBuilderImpl(uCont)
        try {
            builder(scope)
        } catch (e: Throwable) {
            scope.handleBuilderException(e)
        }
        scope.getResult()
    }

suspendCoroutineUninterceptedOrReturn()函数看不到源码,是由编译器生成源码,其作用是:先获取当前挂起函数的续体对象,然后要么挂起当前运行的协程,要么直接返回结果。当在其参数block中返回COROUTINE_SUSPENDED值时,即表明该挂起函数会挂起且不会立即返回任何结果。在block代码块中,先创建一个SelectBuilderImpl对象,然后执行builder,这里的builder是外面传入的,即会执行fizz.onReceive,onReceive是SelectClause1类型,会先执行其get方法,再执行其invoke()方法:

//AbstractChannel
final override val onReceive: SelectClause1<E>
        get() = object : SelectClause1<E> {
            @Suppress("UNCHECKED_CAST")
            override fun <R> registerSelectClause1(select: SelectInstance<R>, block: suspend (E) -> R) {
                registerSelectReceiveMode(select, RECEIVE_THROWS_ON_CLOSE, block as suspend (Any?) -> R)
            }
        }
        
//class SelectBuilderImpl
override fun <Q> SelectClause1<Q>.invoke(block: suspend (Q) -> R) {
		//会执行get()方法返回的SelectClause1的registerSelectClause1()方法
        registerSelectClause1(this@SelectBuilderImpl, block)
    }

//AbstractChannel
private fun <R> registerSelectReceiveMode(select: SelectInstance<R>, receiveMode: Int, block: suspend (Any?) -> R) {
        while (true) {
            if (select.isSelected) return
            if (isEmptyImpl) {
            	//第一次Channel还没有Send元素
                if (enqueueReceiveSelect(select, block, receiveMode)) return
            } else {
                val pollResult = pollSelectInternal(select)
                when {
                    pollResult === ALREADY_SELECTED -> return
                    pollResult === POLL_FAILED -> {} // retry
                    pollResult === RETRY_ATOMIC -> {} // retry
                    else -> block.tryStartBlockUnintercepted(select, receiveMode, pollResult)
                }
            }
        }
    }

registerSelectReceiveMode()函数中,当Channel中没有Send元素时且缓存为空时isEmptyImpl为真,接着会执行enqueueReceiveSelect方法向该Channel的队列中添加一个ReceiveSelect元素:

//AbstractChannel
private fun <R> enqueueReceiveSelect(
    select: SelectInstance<R>,
    block: suspend (Any?) -> R,
    receiveMode: Int
): Boolean {
    val node = ReceiveSelect(this, select, block, receiveMode)
    val result = enqueueReceive(node)
    if (result) select.disposeOnSelect(node)
    return result
}

先用this和block等构造一个ReceiveSelect对象,this即是当前fizz Channel,而block即是fizz.onReceive {}中的lambda表达式。然后调用enqueueReceive()方法将该节点node添加到Channel的queue队列中。最后当添加成功后会调用select的disposeOnSelect()方法使得当该类型为SelectBuilderImpl的select对象被设置为select状态时对node进行一些操作,node继承于DisposableHandle。

//class SelectBuilderImpl
 override fun disposeOnSelect(handle: DisposableHandle) {
        val node = DisposeNode(handle)
        // check-add-check pattern is Ok here since handle.dispose() is safe to be called multiple times
        if (!isSelected) {
            addLast(node) // add handle to list
            // double-check node after adding
            if (!isSelected) return // all ok - still not selected
        }
        // already selected
        handle.dispose()
    }

SelectBuilderImpl继承于LockFreeLinkedListHead,将该ReceiveSelect对象封装成一个DisposeNode对象,然后添加在链表的尾部。
总结:通过对select()函数的代码跟踪,发现其主要做了两件事:

  1. 构造一个ReceiveSelect节点,并添加到ReceiveChannel的queue队列中;
  2. 将ReceiveSelect进一步封装成DisposeNode节点,添加到SelectBuilderImpl的尾部

1.2. send()函数细节

当fizz Channel延迟400ms后,会向通道中发射一个"Fizz"字符串。这时会先从该Channel中取一个元素:

protected open fun offerInternal(element: E): Any {
        while (true) {
            val receive = takeFirstReceiveOrPeekClosed() ?: return OFFER_FAILED
            val token = receive.tryResumeReceive(element, null)
            if (token != null) {
                assert { token === RESUME_TOKEN }
                receive.completeResumeReceive(element)
                return receive.offerResult
            }
        }
    }

这里在fizz Channel中取到的元素receive是上文中添加的ReceiveSelect元素,然后调用其tryResumeReceive()方法:

//class ReceiveSelect
override fun tryResumeReceive(value: E, otherOp: PrepareOp?): Symbol? =
    select.trySelectOther(otherOp) as Symbol?

//class SelectBuilderImpl
override fun trySelectOther(otherOp: PrepareOp?): Any? {
        _state.loop { state -> // lock-free loop on state
            when {
                // Found initial state (not selected yet) -- try to make it selected
                state === this -> {
                    if (otherOp == null) {
                        // regular trySelect -- just mark as select
                        if (!_state.compareAndSet(this, null)) return@loop
                    } else //省略代码
                    doAfterSelect()
                    return RESUME_TOKEN
                }
			//。。。省略代码
            }
        }
    }

//class SelectBuilderImpl
private fun doAfterSelect() {
        parentHandle?.dispose()
        forEach<DisposeNode> {
            it.handle.dispose()
        }
    }

这里将_state设置为null,代表设置select状态,然后执行 doAfterSelect()方法将对链表的其他handler进行处理。这里的handle即是ReceiveSelect对象,执行其dispose()方法:

//class ReceiveSelect
override fun dispose() { // invoked on select completion
            if (remove())
                channel.onReceiveDequeued() // notify cancellation of receive
        }

回到offerInternal()函数中,此时token = RESUME_TOKEN,因此会执行ReceiveSelect的completeResumeReceive()方法:

//class ReceiveSelect
override fun completeResumeReceive(value: E) {
    block.startCoroutine(if (receiveMode == RECEIVE_RESULT) ValueOrClosed.value(value) else value, select.completion)
}

这里的block即是最开始izz.onReceive {}中传入的lambda表示式,调用block的startCoroutine()方法开启一个协程,因为block中可能会调用挂起函数。下面就会执行block中的代码了,因为没有调用挂起函数,所以就打印一个字符串:“Fizz”。到这里就完成了多个Channel中多路复用时一轮数据从发射到接收过程的分析。

1.3 取消或完成时select协程实现细节

在1.1节中,select()函数开启一个新的协程时最后会调用SelectBuilderImpl的getResult()方法,在该方法中实现了取消或完成时的钩子:

//class SelectBuilderImpl
internal fun getResult(): Any? {
        if (!isSelected) initCancellability()
       //省略
    }

//class SelectBuilderImpl
private fun initCancellability() {
        val parent = context[Job] ?: return
        val newRegistration = parent.invokeOnCompletion(
            onCancelling = true, handler = SelectOnCancelling(parent).asHandler)
        parentHandle = newRegistration
        // now check our state _after_ registering
        if (isSelected) newRegistration.dispose()
    }

这里的parent即为最外程的主协程BlockingCoroutine,然后向BlockingCoroutine协程的state变量中添加一个SelectOnCancelling节点,并将该节点赋值给parentHandle成员变量,用于取消或完成时将该节点从BlockingCoroutine协程中移除,通过1.2节可知移除操作是在doAfterSelect()函数中完成的,该函数会继续调用parentHandle的dispose()方法,该方法在其父类JobNode中:

//class JobNode
override fun dispose() = (job as JobSupport).removeNode(this)

这里的job即为BlockingCoroutine协程,这里即将parentHandle从BlockingCoroutine协程的state中移除。

2. Coroutine的多路复用

先看一个Coroutine的多路复用的例子,先定义多个协程:

fun CoroutineScope.asyncString(time: Int) = async {
    delay(time.toLong())
    val s = "Waited for $time ms"
    println("S = $s")
    s
}

fun CoroutineScope.asyncStringsList(): List<Deferred<String>> {
    val random = Random(3)
    return List(2) { asyncString(random.nextInt(1000)) }
}

然后对上面的协程列表进行多路复用:

fun selectThree() {
    runBlocking {
        val list = asyncStringsList()
        val result = select<String> {
            list.withIndex().forEach { (index, deferred) ->
                deferred.onAwait { answer ->
                    "Deferred $index produced answer '$answer'"
                }
            }
        }
        println(result)
        val countActive = list.count { it.isActive }
        println("$countActive coroutines are still active")
    }
}
打印结果为:
S = Waited for 200 ms
Deferred 0 produced answer 'Waited for 200 ms'
1 coroutines are still active
S = Waited for 355 ms

在该例子中,select()函数以及select协程取消和完成时的实现和1.1节、1.3节是一致的,不同的是这里复用的是多个DeferredCoroutine,且采用onAwait 子句放在select中。

2.1. onAwait()子句实现细节

CoroutineScope的async()函数中,根据协程启动方式构造一个DeferredCoroutine或者LazyDeferredCoroutine对象,LazyDeferredCoroutine是DeferredCoroutine的子类,因此这里会先执行onAwait的get()属性方法返回一个SelectClause1实现类,然后再调用SelectClause1的invoke()方法:

//class DeferredCoroutine
override val onAwait: SelectClause1<T> get() = this

//class DeferredCoroutine
override fun <R> registerSelectClause1(select: SelectInstance<R>, block: suspend (T) -> R) =
    registerSelectClause1Internal(select, block)

//class JobSupport
internal fun <T, R> registerSelectClause1Internal(select: SelectInstance<R>, block: suspend (T) -> R) {
        // fast-path -- check state and select/return if needed
        loopOnState { state ->
            if (select.isSelected) return
            if (state !is Incomplete) {
                // already complete -- select result,一般不会走这里路径
                if (select.trySelect()) {
                    if (state is CompletedExceptionally) {
                        select.resumeSelectWithException(state.cause)
                    }
                    else {
                        block.startCoroutineUnintercepted(state.unboxState() as T, select.completion)
                    }
                }
                return
            }
            if (startInternal(state) == 0) {
                // slow-path -- register waiter for completion
                select.disposeOnSelect(invokeOnCompletion(handler = SelectAwaitOnCompletion(this, select, block).asHandler))
                return
            }
        }
    }

当前DeferredCoroutine是继承于JobSupport的,先创建一个继承于JobNode的SelectAwaitOnCompletion对象,然后会调用JobSupport的invokeOnCompletion()方法会在JobSupport的state指向的list中添加一个JobNode节点,invokeOnCompletion()方法返回的也是该SelectAwaitOnCompletion节点。接着会调用SelectBuilderImpl的disposeOnSelect()方法继续将该节点添加到SelectBuilderImpl末尾,SelectBuilderImpl继承于LockFreeLinkedListHead。

2.2. 多路复用的多个协程中某个协程恢复细节

在本节的例程中select多路复用了两个DeferredCoroutine协程,每个DeferredCoroutine协程的state对象中又分别添加了一个SelectAwaitOnCompletion节点,当某个协程恢复时先会将该DeferredCoroutine协程从其父协程BlockingCoroutine中删除,然后会循环调用该DeferredCoroutine协程state指向的list中每个节点的invoke()方法。

BaseContinuationImpl AbstractCoroutine JobSupport resumeWith() tryMakeCompleting() tryMakeCompletingSlowPath() completeStateFinalization() notifyHandlers() BaseContinuationImpl AbstractCoroutine JobSupport

在JobSupport的notifyHandlers()方法中会循环调用state中list的每个node的invoke()方法:

//class SelectAwaitOnCompletion
override fun invoke(cause: Throwable?) {
        if (select.trySelect())
            job.selectAwaitCompletion(select, block)
    }

这里SelectBuilderImpl的trySelect()方法进而会调用trySelectOther()方法来回调处理添加在SelectBuilderImpl中的所有SelectAwaitOnCompletion节点,并调用其dispose()方法,这里SelectAwaitOnCompletion没有重写dispose()方法,因此会调用其父类JobNode的dispose()方法:

//class JobNode
override fun dispose() = (job as JobSupport).removeNode(this)

这里的job即是DeferredCoroutine协程,其内部的删除逻辑同样是在state上处理的:

internal fun removeNode(node: JobNode<*>) {
        loopOnState { state ->
            when (state) {
            	//如果当前state指向一个JobNode对象,则将该state设置为EMPTY_ACTIVE
                is JobNode<*> -> { // SINGE/SINGLE+ state -- one completion handler
                    if (state !== node) return // a different job node --> we were already removed
                    // try remove and revert back to empty state
                    if (_state.compareAndSet(state, EMPTY_ACTIVE)) return
                }
                //如果当前state指向一个Incomplete对象,则将该node从list链表中删除
                is Incomplete -> { // may have a list of completion handlers
                    // remove node from the list if there is a list
                    if (state.list != null) node.remove()
                    return
                }
                else -> return // it is complete and does not have any completion handlers
            }
        }
    }

回到SelectAwaitOnCompletion中的invoke()方法,这里job是DeferredCoroutine协程,block即是onAwait{}中传入的lambda表达式,trySelect()返回true,因此会执行其selectAwaitCompletion()方法:

//class JobSupport
internal fun <T, R> selectAwaitCompletion(select: SelectInstance<R>, block: suspend (T) -> R) {
        val state = this.state
        // Note: await is non-atomic (can be cancelled while dispatched)
        if (state is CompletedExceptionally)
            select.resumeSelectWithException(state.cause)
        else
            block.startCoroutineCancellable(state.unboxState() as T, select.completion)
    }

这里的state为“Waited for 200 ms”,因此调用block的startCoroutineCancellable()方法后会马上执行block中的代码,因此即会打印出:Deferred 0 produced answer ‘Waited for 200 ms’。到这里一个协程执行完毕产生一个"Waited for 200 ms"字符串到执行onAwait{}代码块打印“Deferred 0 produced answer ‘Waited for 200 ms’”的过程即分析完毕了。那还有一个问题:等待另一个协程执行完毕也产生一个字符串“Waited for xxx ms”后为什么不再执行onAwait{}代码块了呢?答案是添加在DeferredCoroutine协程中的SelectAwaitOnCompletion已经被删除了,当另一个协程执行完毕后已删除的SelectAwaitOnCompletion对象的invoke()方法不会被调用,因此也不会执行2.2节过程的代码,因此onAwait{}代码块不会再执行了。

Logo

开源鸿蒙跨平台开发社区汇聚开发者与厂商,共建“一次开发,多端部署”的开源生态,致力于降低跨端开发门槛,推动万物智联创新。

更多推荐