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

Spring Boot集成Cloud PubSub重复发消息问题求助

问题原因分析
  1. 通道复用引发消息循环与重复处理
    你共用了inputMessageChannel通道,导致以下问题:

    • 入站适配器从chat-sub订阅接收消息后,发送到该通道
    • 通道同时被messageReceiver(处理消息)和messageSender(转发到chat主题)监听
    • 形成消息循环:订阅消息→进入通道→转发回主题→订阅再次接收,无限重复。同时网关发送的消息进入通道后,messageReceiver因找不到gcp_pubsub_original_message头抛出异常,触发框架重试,进一步加剧重复。
  2. 手动ACK与自动ACK冲突
    你配置了ackMode = AckMode.AUTO,框架会自动完成消息ACK,但messageReceiver中又手动调用message.ack(),导致同一个ACK ID被重复提交,PubSub返回无效ACK ID异常。

  3. 缺失头部触发重试
    网关发送的消息没有gcp_pubsub_original_message头,但messageReceiver强制要求该参数,抛出MessageHandlingException。Spring Integration默认会重试处理失败的消息,导致消息被反复处理,产生大量重复。

解决步骤

1. 拆分通道,阻断循环与冲突

将入站消息处理和出站消息发送的通道分开,避免共用:

// 入站专用通道:处理订阅接收的消息
@Bean
fun inboundProcessingChannel(): MessageChannel {
    return DirectChannel()
}

// 出站专用通道:发送消息到PubSub主题
@Bean
fun outboundPublishChannel(): MessageChannel {
    return DirectChannel()
}

// 入站适配器:将chat-sub消息发送到入站通道
@Bean
fun inboundChannelAdapter(
    @Qualifier("inboundProcessingChannel") messageChannel: MessageChannel,
    pubSubTemplate: PubSubTemplate
): PubSubInboundChannelAdapter {
    val adapter = PubSubInboundChannelAdapter(pubSubTemplate, "chat-sub")
    adapter.outputChannel = messageChannel
    adapter.ackMode = AckMode.AUTO
    adapter.payloadType = String::class.java
    return adapter
}

// 入站消息处理器:监听入站通道
@ServiceActivator(inputChannel = "inboundProcessingChannel")
fun messageReceiver(
    payload: String,
    @Header(GcpPubSubHeaders.ORIGINAL_MESSAGE) message: BasicAcknowledgeablePubsubMessage
) {
    val model = gson.fromJson(payload, ChatMessageToPubSub::class.java)
    messagingTemplate.convertAndSend("/sub/message/" + model.roomId.toString(), ChatMessageToClient(model.userId, model.message))
    // 注意:AUTO模式下删除手动ack(),避免重复提交
}

// 出站适配器:监听出站通道,发送到chat主题
@Bean
@ServiceActivator(inputChannel = "outboundPublishChannel")
fun messageSender(pubsubTemplate: PubSubTemplate): MessageHandler {
    val adapter = PubSubMessageHandler(pubsubTemplate, "chat")
    adapter.setSuccessCallback { ackId, message -> println("send success : " + message.payload) }
    adapter.setFailureCallback { cause, message -> println("send fail : " + message.payload) }
    return adapter
}

// 网关:将消息发送到出站通道
@MessagingGateway(defaultRequestChannel = "outboundPublishChannel")
interface PubSubOutBoundGateway {
    @Throws(MessagingException::class)
    fun sendToPubsub(text : String)
}

2. 统一ACK策略

  • 若使用AckMode.AUTO,务必删除messageReceiver中的message.ack()调用,由框架自动处理ACK
  • 若需手动控制ACK,将ackMode改为AckMode.MANUAL,确保只调用一次ack()或nack(),禁止重复操作

3. 禁用不必要的重试

如果需要全局控制重试行为,可在出站适配器中配置禁用重试:

// 在messageSender Bean中添加
adapter.setRetryTemplate(RetryTemplate().apply {
    setRetryPolicy(NeverRetryPolicy())
})

4. 清理积压消息

重启服务前,通过GCP控制台手动确认chat-sub订阅中的未处理消息,避免重启后再次触发重复推送;客户端可添加基于消息ID的去重逻辑,防止重复处理。

额外说明

BasicAcknowledgeablePubsubMessage是Spring Cloud GCP内部实现的接口(实现类为AcknowledgeablePubsubMessage),无需手动实例化,只有从PubSub订阅接收的消息才会携带该头部,网关发送的消息不存在此头部,这也是之前缺失头部异常的根源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 15:38:08