Spring WebFlux handler for kotlin SharedFlow

Viewed 40

I can see the following example working in Spring WebFlux handler for a flow builder:

suspend fun getDummyFlow(req: ServerRequest): ServerResponse {
    val flow = flow<String> { // flow builder
        for (i in 1..3) {
            delay(1000) // pretend we are doing something useful here
            emit("<p>Hello $i</p>") // emit next value
        }
    }
    return ServerResponse
        .ok()
        .contentType(MediaType.TEXT_HTML)
        .bodyAndAwait(flow)
}

Yet, I need to build a flow with a MutableSharedFlow which is not working in Spring Web Flux. Here it is an example:

suspend fun getDummyFlow(req: ServerRequest): ServerResponse {
    return coroutineScope {
        val flow = MutableSharedFlow<String>()
        launch {
            for (i in 1..3) {
                delay(1000) // pretend we are doing something useful here
                flow.emit("<p>Hello $i</p>") // emit next value
            }
        }
        ServerResponse
            .ok()
            .contentType(MediaType.TEXT_HTML)
            .bodyAndAwait(
                flow
                    .asSharedFlow()
                    .take(3)
            )
    }

My implementation is based on the example of SharedFlow documentation.

Yet, any HTTP GET request to this endpoint stays pending and waiting for a response, whereas the former example with flow builder receives the response progressively and fine.

I have already traced my code in debug and I see .bodyAndAwait(..) being called and then emit() in both cases.

0 Answers
Related