You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()

关键说明

  1. 串行处理队列:用ConcurrentHashMap存储每个排序键的处理Job,新消息会等待上一条同键消息的Job执行完成后再启动,确保同键消息串行处理。
  2. 非阻塞延迟:使用Kotlin协程的delay()实现非阻塞等待,不会阻塞线程,相比Thread.sleep性能更优。
  3. FlowControl限制:通过setMaxOutstandingElementCount(1L)限制订阅者本地未处理消息数量,结合排序键的顺序投递特性,避免Pub/Sub一次性推送多条同键消息到本地。
  4. 异常处理:捕获处理异常,调用nack()让Pub/Sub重新投递消息,保证消费可靠性。

效果验证

发送5条同排序键消息后,每条消息会间隔10秒依次打印,实现逐条延迟消费的需求。

内容的提问来源于stack exchange,提问作者Mark

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.10 21:55:19