Coroutine之Channel的多路复用原理浅析
目录
根据 官方文档可知select 表达式可以同时等待多个挂起函数,并选择第一个可用的。select具有一下功能可以细分如下:
- 在多个Channel中选择第一个到达的进行处理;
- 可以对主-从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()函数的代码跟踪,发现其主要做了两件事:
- 构造一个ReceiveSelect节点,并添加到ReceiveChannel的queue队列中;
- 将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()方法。
在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{}代码块不会再执行了。
更多推荐


所有评论(0)