How can we use map instead of flatMapMerge in Kotlin coroutines?

Viewed 249

I'm migrating a Webflux codebase that uses Reactor types (Mono and Flux) to Kotlin coroutines.

Consider this excerpt with Flux:


fun getClientProjects(clientId: String) : Flux<Project> {
    return projectReactiveRepository.findByClientId(clientId)
}

fun getAllProjects() : Flux<Project> {
    return clientRepository.findAll().flatMap { client -> getProjects(client.id) }
}

It gets all clients and for each of them invokes a reactive repository concurrently, subscribing eagerly to the nested Fluxes.

This could be translated to coroutines as shown below:

suspend fun getClientProjects(clientId: String) : Flow<Project> {
    return projectReactiveRepository.findByClientId(clientId).asFlow()
}

fun getAllProjects() : Flow<Project> {
    return clientRepository.findAll().asFlow().flatMapMerge { client -> getProjects(client.id) }
}

Flow.flatMapMerge javadoc states:

Note that even though this operator looks very familiar, we discourage its usage in a regular application-specific flows. Most likely, suspending operation in map operator will be sufficient and linear transformations are much easier to reason about.

How could I replace flatMapMerge with map and keep the concurrency of the nested Flows?

0 Answers
Related