Google Pub/Sub同排序键消息延迟消费问题(非阻塞实现)
Google Pub/Sub 同排序键消息间隔延迟消费问题解决
问题描述
已成功订阅Google Pub/Sub消息,主题配置了排序键,期望实现:每个排序键每次仅消费1条消息,延迟后再消费该排序键的下一条消息。
当前代码
import com.google.cloud.pubsub.v1.MessageReceiver import com.google.cloud.pubsub.v1.Subscriber val receiver = MessageReceiver { message, consumer -> GlobalScope.launch { delay(10 * 1000L) println(message.data.toStringUtf8()) consumer.ack() } } subscriber = Subscriber.newBuilder(subscriptionName, receiver).build() subscriber!!.startAsync().awaitRunning()
当前异常表现
一次性发送5条带相同排序键的消息后,60秒后打印第一条消息,再过60秒一次性打印剩余4条消息,不符合逐条间隔延迟的需求。
需要解决两个问题:
- 如何修改代码实现同排序键消息逐条间隔延迟消费?
- 能否用非阻塞延迟替代
Thread.sleep解决问题?
解决方案
核心问题分析
Pub/Sub的排序键机制会按顺序投递同键消息,但默认订阅者会预取多条同键消息到本地队列。当前代码中每条消息都直接启动协程计时,订阅者会持续推送同键消息到Receiver,导致所有协程几乎同时开始计时;加上Pub/Sub的默认重试与预取机制,最终出现批量打印的情况。
修改思路
为每个排序键维护串行处理队列,确保同键消息仅在前一条处理完成(含延迟)并ACK后,再启动下一条的处理;同时用Kotlin协程的非阻塞延迟实现等待,完全不需要Thread.sleep。
修改后的代码
import com.google.cloud.pubsub.v1.MessageReceiver import com.google.cloud.pubsub.v1.Subscriber import kotlinx.coroutines.* import java.util.concurrent.ConcurrentHashMap // 为每个排序键维护串行处理的Job,保证同键消息顺序执行 private val keyProcessingJobs = ConcurrentHashMap<String, Job>() // 自定义协程作用域,避免使用GlobalScope带来的生命周期管理问题 private val coroutineScope = CoroutineScope(Dispatchers.Default + SupervisorJob()) val receiver = MessageReceiver { message, consumer -> val orderingKey = message.orderingKey // 处理空排序键的消息(可选:按默认逻辑或跳过) if (orderingKey.isNullOrEmpty()) { coroutineScope.launch { delay(10 * 1000L) println(message.data.toStringUtf8()) consumer.ack() } return@MessageReceiver } // 为当前排序键构建串行处理链 keyProcessingJobs.compute(orderingKey) { _, existingJob -> coroutineScope.launch { // 等待上一条同键消息处理完成 existingJob?.join() try { // 非阻塞延迟10秒 delay(10 * 1000L) println(message.data.toStringUtf8()) consumer.ack() } catch (e: Exception) { // 异常时NACK,让Pub/Sub重新投递消息 consumer.nack() } } } } subscriber = Subscriber.newBuilder(subscriptionName, receiver) // 关键配置:限制本地未处理消息数量为1,避免Pub/Sub批量推送同键消息 .setFlowControlSettings( Subscriber.Builder.FlowControlSettings.newBuilder() .setMaxOutstandingElementCount(1L) .build() ) .build() subscriber!!.startAsync().awaitRunning()
关键说明
- 串行处理队列:用
ConcurrentHashMap存储每个排序键的处理Job,新消息会等待上一条同键消息的Job执行完成后再启动,确保同键消息串行处理。 - 非阻塞延迟:使用Kotlin协程的
delay()实现非阻塞等待,不会阻塞线程,相比Thread.sleep性能更优。 - FlowControl限制:通过
setMaxOutstandingElementCount(1L)限制订阅者本地未处理消息数量,结合排序键的顺序投递特性,避免Pub/Sub一次性推送多条同键消息到本地。 - 异常处理:捕获处理异常,调用
nack()让Pub/Sub重新投递消息,保证消费可靠性。
效果验证
发送5条同排序键消息后,每条消息会间隔10秒依次打印,实现逐条延迟消费的需求。
内容的提问来源于stack exchange,提问作者Mark
相关产品推荐
相关产品推荐

