【问题标题】:Coroutine: Deferred operations in a List run sequentially.协程:List 中的延迟操作按顺序运行。
【发布时间】:2019-04-13 18:33:11
【问题描述】:

我有一个List 的参数用于执行下载。 我将该列表的元素映射到执行下载的Deferred;然后,列表的forEach 元素,我调用await,但显然下载是按顺序执行的。

这是我的功能:

suspend fun syncFiles() = coroutineScope {
    remoteRepository.requiredFiles()
        .filter { localRepository.needToDownload( it.name, it.md5 ) }
        .map { async { downloader( it ) } }
        .forEach { deferredResult ->
​
            when ( val result = deferredResult.await() ) {
                is DownloadResult.Layout ->  localRepository.storeLayout( result.content )
                is DownloadResult.StringR -> localRepository.storeFile( result )
            }
        }
}

这是我的测试:

private val useCase = SyncUseCaseImpl.Factory(
        mockk { // downloader
            coEvery { this@mockk.invoke( any() ) } coAnswers { delay(1000 );any() }
        },
        ...
    ).newInstance()
​
@Test
fun `syncFiles downloadConcurrently`() = runBlocking {
    val requiredFilesCount = useCase.remoteRepository.requiredFiles().size
    assert( requiredFilesCount ).isEqualTo( 3 )
​
    val time = measureTimeMillis {
        useCase.syncFiles()
    }
​
    assert( time ).isBetween( 1000, 1100 )
}

这是我的结果:expected to be between:<1000L> and <1100L> but was:<3081L>

我觉得很奇怪,因为这两个虚拟测试正确完成,也许我遗漏了什么(?)

@Test // OK
fun test() = runBlocking {
    val a = async { delay(1000 ) }
    val b = async { delay(1000 ) }
    val c = async { delay(1000 ) } ​
    val time = measureTimeMillis {
        a.await()
        b.await()
        c.await()
    } ​
    assert( time ).isBetween( 1000, 1100 )
} ​

@Test // OK
fun test() = runBlocking {
    val wasteTime: suspend () -> Unit = { delay(1000 ) }
    suspend fun wasteTimeConcurrently() = listOf( wasteTime, wasteTime, wasteTime )
            .map { async { it() } }
            .forEach { it.await() } ​
    val time = measureTimeMillis {
        wasteTimeConcurrently()
    } ​
    assert( time ).isBetween( 1000, 1100 )
}

【问题讨论】:

  • 您的代码看起来正确。您确定您的downloader 同时运行并且有多个连接吗?另一个原因可能是你的localRepository,也许它只按顺序运行?
  • 一切都被嘲笑了,下载器在延迟 1 秒后返回 any()
  • 首先要怀疑的是模拟实现。也许它会懒惰地启动并仅在您调用await时执行delay
  • 嗯,那会很奇怪????但我不是说不可能,我会调试它,谢谢

标签: concurrency kotlin coroutine kotlinx.coroutines


【解决方案1】:

问题出在mockk

如果您查看coAnswer 函数的代码,您会发现(API.kt + InternalPlatformDsl.kt):

infix fun coAnswers(answer: suspend MockKAnswerScope<T, B>.(Call) -> T) = answers {
    InternalPlatformDsl.runCoroutine {
        answer(it)
    }
}

runCoroutine 看起来像这样。

actual fun <T> runCoroutine(block: suspend () -> T): T {
    return runBlocking {
        block()
    }
}

如你所见,coAnswer 是一个非挂起函数,它使用runBlocking 启动一个新的协程。

让我们看一个例子:

val mock =  mockk<Downloader> {
    coEvery {
        this@mockk.download()
    } coAnswers {
        delay(1000)
    }
}

val a = async {
    mock.download()
}

mockk 执行coAnswer-block (delay()) 时,它会启动一个人工协程范围,执行给定的块并等待(阻塞当前线程:runBlocking)直到该块完成。所以答案块只有在delay(1000) 完成后才会返回。

表示从coAnswer运行的所有协程都按顺序执行。

【讨论】:

  • 是的,但我认为这很正常;但总执行时间应该比一秒多一点,因为它们可以几乎一起开始。我错了吗? ?
  • 我的测试用例是错误的,在这种情况下你是对的 :-) 但问题还是一样。我会改变我的答案。
  • 如果您仍在使用代码,请尝试measureTimeMillis { a.await(); b.await(); c.await() } 吗?我的电脑在我的旅行包里,我会在 2 小时内醒来,我无法入睡??
  • 我刚刚改了答案。希望解释足够好。
  • 不是真的,说实话?但我对协程很陌生,所以我想是我的错。明天我会检查一下,现在我想我可以停止思考了。谢谢?
【解决方案2】:

如果作业阻塞了整个线程,则可能会发生这种情况,例如 IO 绑定任务阻塞了整个线程的执行,从而阻塞了该线程上的所有其他协程。如果您使用 Kotlin JVM,请尝试调用 async(IO) { } 在 IO 调度程序下运行 couroutine,以便 couroutine 环境现在知道该作业将阻塞整个线程并相应地运行。

在这里查看其他调度员:https://kotlinlang.org/docs/reference/coroutines/coroutine-context-and-dispatchers.html#dispatchers-and-threads

【讨论】:

    猜你喜欢
    • 2016-03-24
    • 1970-01-01
    • 2015-07-24
    • 2014-09-15
    • 2015-11-04
    • 2013-04-29
    • 1970-01-01
    • 2013-01-10
    • 2013-11-19
    相关资源
    最近更新 更多