Is is possible to bridge between reactive code and kotlin coroutines without blocking with runBlocking?

Viewed 651

I writing a KafkaConsumer in Kotlin using the reactive framework, the problem is the whole project structure is based on Kotlin coroutines and now the kafka consumer follows a Flux publisher pipeline,

I got it to work with runBlocking however I am aware this is not a good idea to have blocking code in our project,

I tried using @KafkaListener (fails when adding suspend modifier),

import com.github.avrokotlin.avro4k.Avro
import kotlinx.coroutines.runBlocking
import org.apache.avro.generic.GenericRecord
import org.apache.kafka.clients.consumer.ConsumerRecord
import org.springframework.boot.CommandLineRunner
import org.springframework.kafka.core.reactive.ReactiveKafkaConsumerTemplate
import org.springframework.stereotype.Component

@Component
class KafkaConsumer(
    val slackNotificationService: SlackNotificationService,
    val consumerTemplate: ReactiveKafkaConsumerTemplate<String, GenericRecord>
) : CommandLineRunner {
    suspend fun sendNotification(record: ConsumerRecord<String, GenericRecord>) {
        val tagNotification = Avro.default.fromRecord(TagNotification.serializer(), record.value())
        slackNotificationService.notifyUsers(tagNotification)
    }

    override fun run(vararg args: String?) {
        consumerTemplate
            .receiveAutoAck()
            .subscribe {
                runBlocking {
                    sendNotification(it)
                }
            }
    }
}

I can successfully receive the kafka messages and all the rest of the project is working fine, but I couldn't figure out to create this non blocking bridge here,

Does anyone know if there's a better way to handle this?

Thank you in advance :)

1 Answers

If you want to invoke sendNotification() asynchronously then create a CoroutineScope and launch coroutines with it:

class KafkaConsumer(
    ...
    private val coroutineScope = CoroutineScope(Dispatchers.Default)
    ...
            .subscribe {
                coroutineScope.launch { 
                    sendNotification(it)
                }
            }

If KafkaConsumer may be destroyed/shutdown then it is advised to invoke coroutineScope.cancel() when it happens.

Related