如何在不使用runBlocking阻塞的前提下桥接响应式代码与Kotlin协程?
解决方案
方案1:基于Reactor与协程桥接适配现有响应式消费者
这是对现有代码改动最小的非阻塞方案,依赖kotlinx-coroutines-reactor库实现Reactor Publisher到Kotlin协程Flow的转换,全程无阻塞调用。
第一步:添加依赖
Maven配置:
<dependency> <groupId>org.jetbrains.kotlinx</groupId> <artifactId>kotlinx-coroutines-reactor</artifactId> <version>匹配当前项目Kotlin版本即可</version> </dependency>
Gradle配置:
implementation "org.jetbrains.kotlinx:kotlinx-coroutines-reactor:匹配当前项目Kotlin版本"
第二步:修改消费者代码
给消费者组件绑定独立的协程作用域,实现生命周期管理,避免内存泄漏:
import com.github.avrokotlin.avro4k.Avro import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.flow.collect import kotlinx.coroutines.launch import kotlinx.coroutines.reactor.asFlow 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 import jakarta.annotation.PreDestroy @Component class KafkaConsumer( val slackNotificationService: SlackNotificationService, val consumerTemplate: ReactiveKafkaConsumerTemplate<String, GenericRecord> ) : CommandLineRunner { // 绑定组件生命周期的协程作用域,SupervisorJob保证单条消息处理失败不影响整体消费 private val consumerScope = CoroutineScope(Dispatchers.IO + SupervisorJob()) 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?) { consumerScope.launch { consumerTemplate .receiveAutoAck() .asFlow() // 将Flux转为协程Flow,无缝衔接非阻塞处理 .collect { // 直接调用suspend方法,无需阻塞 runCatching { sendNotification(it) }.onFailure { e -> // 此处可添加自定义异常处理逻辑,避免单条消息异常终止整个消费流 e.printStackTrace() } } } } // 组件销毁时取消协程作用域,避免资源泄漏 @PreDestroy fun stopConsumer() { consumerScope.cancel() } }
方案2:使用@KafkaListener原生协程支持
如果使用Spring Boot 3.0+ / Spring Kafka 2.8+ 版本,@KafkaListener已经原生支持suspend修饰的方法,无需借助响应式模板类,代码更简洁:
import com.github.avrokotlin.avro4k.Avro import org.apache.avro.generic.GenericRecord import org.apache.kafka.clients.consumer.ConsumerRecord import org.springframework.kafka.annotation.KafkaListener import org.springframework.stereotype.Component @Component class KafkaConsumer( val slackNotificationService: SlackNotificationService ) { suspend fun sendNotification(record: ConsumerRecord<String, GenericRecord>) { val tagNotification = Avro.default.fromRecord(TagNotification.serializer(), record.value()) slackNotificationService.notifyUsers(tagNotification) } @KafkaListener(topics = ["你的Topic名称"], groupId = "你的消费者组ID") suspend fun consume(record: ConsumerRecord<String, GenericRecord>) { sendNotification(record) } }
之前添加suspend修饰符运行失败大概率是Spring版本过低,或者没有配置对应Avro消息转换器导致参数解析失败,升级版本并配置好GenericMessageConverter即可正常使用。
原有方案问题说明
之前使用runBlocking会阻塞Reactor调度线程,会直接影响整个应用的响应式事件循环运行,高并发场景下会出现严重的消费卡顿、线程耗尽问题,上述两种方案都完全规避了阻塞操作,符合非阻塞架构设计。
内容的提问来源于stack exchange,提问作者Lucas
相关产品推荐
相关产品推荐

