Spring Boot集成Cloud PubSub重复发消息问题求助
问题原因分析
通道复用引发消息循环与重复处理
你共用了inputMessageChannel通道,导致以下问题:- 入站适配器从
chat-sub订阅接收消息后,发送到该通道 - 通道同时被
messageReceiver(处理消息)和messageSender(转发到chat主题)监听 - 形成消息循环:订阅消息→进入通道→转发回主题→订阅再次接收,无限重复。同时网关发送的消息进入通道后,
messageReceiver因找不到gcp_pubsub_original_message头抛出异常,触发框架重试,进一步加剧重复。
- 入站适配器从
手动ACK与自动ACK冲突
你配置了ackMode = AckMode.AUTO,框架会自动完成消息ACK,但messageReceiver中又手动调用message.ack(),导致同一个ACK ID被重复提交,PubSub返回无效ACK ID异常。缺失头部触发重试
网关发送的消息没有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
相关产品推荐
相关产品推荐

