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?