Cahining coroutines by using extension functions in Kotlin

Viewed 110

I want to chain 3 coroutines by using Kotlin's extension functions. I know how to do it with regular ones, but can't manage it with extension functions. In fact, in the 2nd coroutine I can receive only one data sent from the 1st coroutine, but that's all. The program works but all I get on the console is Doc: 1st Document. What I'm doing wrong?

fun main(args: Array<String>) = runBlocking {
    produceDocs().docLength().report().consumeEach {
        println(it)
    }
}

private fun CoroutineScope.produceDocs() = produce {
    fun getDocs(): List<String> {
        return listOf("1st Document", "2nd Newer Document")
    }
    while (this.isActive) {
        val docs = getDocs()
        for (doc in docs) {
            send(doc)
        }
        delay(TimeUnit.SECONDS.toMillis(2))
    }
}

private suspend fun ReceiveChannel<String>.docLength(): ReceiveChannel<Int> = coroutineScope {
    val docsChannel: ReceiveChannel<String> = this@docLength

    produce {
        for (doc in docsChannel) {
            println("Doc: $doc") // OK. This works.
            send(doc.count()) // ??? Not sure where this sends data to?
        }
    }
}

private suspend fun ReceiveChannel<Int>.report(): ReceiveChannel<String> = coroutineScope {
    val docLengthChannel: ReceiveChannel<Int> = this@report

    produce {
        for (len in docLengthChannel) {
            println("Length: $len") // !!! Nothing arrived.
            send("Report. Document contains $len characters.")
        }
    }
}
1 Answers

You have to consume each channel independently in order to make emissions go through the chain, otherwise the first emission will never be consumed:

private fun CoroutineScope.produceDocs() = produce {
    fun getDocs(): List<String> {
        return listOf("1st Document", "2nd Newer Document")
    }

    while (this.isActive) {
        val docs = getDocs()
        for (doc in docs) {
            send(doc)
        }
        delay(TimeUnit.SECONDS.toMillis(2))
    }
}

private suspend fun ReceiveChannel<String>.docLength() : ReceiveChannel<Int> = CoroutineScope(coroutineContext).produce {
    for (doc in this@docLength) {
        println("Doc: $doc") // OK. This works.
        send(doc.count()) // ??? Not sure where this sends data to?
    }
}

private suspend fun ReceiveChannel<Int>.report(): ReceiveChannel<String> = CoroutineScope(coroutineContext).produce {
    for (len in this@report) {
        println("Length: $len") // !!! Nothing arrived.
        send("Report. Document contains $len characters.")
    }
}

I suggest you a better approach to do the exact same thing using Flow:

private fun produceDocs(): Flow<String> = flow {
    fun getDocs(): List<String> {
        return listOf("1st Document", "2nd Newer Document")
    }

    while (true) {
        val docs = getDocs()
        for (doc in docs) {
            emit(doc)
        }
        delay(TimeUnit.SECONDS.toMillis(2))
    }
}

private fun Flow<String>.docLength(): Flow<Int> = flow {
    collect { doc ->
        println("Doc: $doc")
        emit(doc.count())
    }
}

private fun Flow<Int>.report(): Flow<String> = flow {
    collect { len ->
        println("Length: $len")
        emit("Report. Document contains $len characters.")
    }
}

Or better like this:

private fun produceDocs(): Flow<String> = flow {
    fun getDocs(): List<String> {
        return listOf("1st Document", "2nd Newer Document")
    }

    while (true) {
        val docs = getDocs()
        for (doc in docs) {
            emit(doc)
        }
        delay(TimeUnit.SECONDS.toMillis(2))
    }
}

private fun Flow<String>.docLength(): Flow<Int> = transform { doc ->
    println("Doc: $doc")
    emit(doc.count())
}

private fun Flow<Int>.report(): Flow<String> = transform { len ->
    println("Length: $len")
    emit("Report. Document contains $len characters.")
}

And collect it like this:

produceDocs().docLength().report().collect {
    println(it)
}

Or even better like this:

produceDocs()
    .map { doc ->
        println("Doc: $doc")
        doc.count()
    }
    .map { len ->
        println("Length: $len")
        "Report. Document contains $len characters."
    }
    .collect {
        println(it)
    }
Related