Spring Boot中PubSubInboundChannelAdapter消息重复问题咨询
问题背景
我正在开发Spring Boot应用,作为Google PubSub订阅者,要求实现Exactly-once delivery(恰好一次投递)。遵循Google Cloud官方指南,但未直接使用PubSub客户端库,而是采用Spring Integration。使用PubSubInboundChannelAdapter时,同ID消息被多次重复消费。
配置代码
PubSub配置类
@Slf4j @Configuration public class PubSubConfig { @Value("${values.gcp.pubsub.subscription.name}") private String subscriptionName; /** * 实现Java对象与JSON的序列化/反序列化,支持Cloud Pub/Sub的JSON消息 payload * * @param objectMapper 使用的对象映射器 * @return Jackson消息转换器 */ @Bean public JacksonPubSubMessageConverter jacksonPubSubMessageConverter(ObjectMapper objectMapper) { return new JacksonPubSubMessageConverter(objectMapper); } @Bean public MessageChannel pubsubInputChannel() { return new DirectChannel(); } @Bean public PubSubInboundChannelAdapter messageChannelAdapter( @Qualifier("pubsubInputChannel") MessageChannel inputChannel, PubSubTemplate pubSubTemplate) { PubSubInboundChannelAdapter adapter = new PubSubInboundChannelAdapter(pubSubTemplate, subscriptionName); adapter.setOutputChannel(inputChannel); adapter.setPayloadType(MyObjectThatNeedBeUnique.class); adapter.setAckMode(AckMode.AUTO_ACK); return adapter; } }
消息监听器代码
@Slf4j @RequiredArgsConstructor(onConstructor = @__(@Autowired)) @Component public class CreateVMListener { @ServiceActivator(inputChannel = "pubsubInputChannel") public void createVMListener(@Payload MyObjectThatNeedBeUnique payload, @Header(GcpPubSubHeaders.ORIGINAL_MESSAGE) BasicAcknowledgeablePubsubMessage message) throws IOException, ExecutionException, InterruptedException, TimeoutException { log.info("Message arrived! Payload: " + payload.toString() + " | MessageId: " + message.getPubsubMessage().getMessageId() ); // 耗时5分钟的业务处理逻辑 } }
配置与日志信息
已在application.yml和Google Cloud PubSub控制台将max-ack-extension-period设置为600秒。
日志显示重复消费问题:
2023-02-05 23:22:20.055 [thread1] INFO - Message arrived! Payload: MyObjectThatNeedBeUnique (userId=432) | MessageID: 6846773022764035 2023-02-05 23:22:31.969 [thread2] INFO - Message arrived! Payload: MyObjectThatNeedBeUnique (userId=432) | MessageID: 6846773022764035 2023-02-05 23:23:33.028 [thread3] INFO - Message arrived! Payload: MyObjectThatNeedBeUnique (userId=432) | MessageID: 6846773022764035 2023-02-05 23:24:34.055 [thread4] INFO - Message arrived! Payload: MyObjectThatNeedBeUnique (userId=432) | MessageID: 6846773022764035
同一条消息在短时间内被4个线程同时处理。
疑问
- 为何出现消息重复?
- 如何在保留
PubSubInboundChannelAdapter高效流式拉取的前提下避免重复?
补充:切换为同步的PubSubMessageSource后问题解决,但该方案性能较低,希望找到使用PubSubInboundChannelAdapter的解决办法。
同步方案代码:
@Bean @InboundChannelAdapter(channel = "pubsubInputChannel") public MessageSource<Object> pubsubAdapter(PubSubTemplate pubSubTemplate) { PubSubMessageSource messageSource = new PubSubMessageSource(pubSubTemplate, createVMSubscriptionName); messageSource.setAckMode(AckMode.AUTO_ACK); messageSource.setPayloadType(CreateOrStartVM.class); messageSource.setBlockOnPull(true); return messageSource; }
解答
问题原因
- AUTO_ACK模式时机错误:
PubSubInboundChannelAdapter采用异步拉取,AUTO_ACK模式下会在消息发送到通道后立即确认,而非等待业务处理完成。如果业务处理耗时超过PubSub的ackDeadline(即使设置了max-ack-extension-period,客户端无主动续期时,PubSub仍会判定消息处理超时),就会触发重新投递。 - 流式拉取的并发特性:
PubSubInboundChannelAdapter基于流式拉取,默认一次性拉取多条消息分发到通道,而DirectChannel默认使用多线程线程池处理消息。若消息提前确认但业务处理超时/失败,PubSub重投后会被线程池的多个线程同时处理,导致重复消费。 - 缺少主动续期机制:
AUTO_ACK模式下,适配器不会自动为耗时较长的业务续期消息的ackDeadline,最终触发PubSub超时重投。
解决方案
要保留PubSubInboundChannelAdapter的流式拉取高效性,同时避免重复消费,可按以下步骤调整:
1. 切换为MANUAL Ack模式,手动管理确认与续期
将适配器的AckMode改为MANUAL,在业务处理完成后手动确认消息,同时在处理过程中主动续期ackDeadline:
@Bean public PubSubInboundChannelAdapter messageChannelAdapter( @Qualifier("pubsubInputChannel") MessageChannel inputChannel, PubSubTemplate pubSubTemplate) { PubSubInboundChannelAdapter adapter = new PubSubInboundChannelAdapter(pubSubTemplate, subscriptionName); adapter.setOutputChannel(inputChannel); adapter.setPayloadType(MyObjectThatNeedBeUnique.class); adapter.setAckMode(AckMode.MANUAL); // 改为手动确认模式 return adapter; }
修改监听器代码,添加手动确认和续期逻辑:
@Slf4j @RequiredArgsConstructor(onConstructor = @__(@Autowired)) @Component public class CreateVMListener { private final PubSubTemplate pubSubTemplate; // 续期间隔建议小于ackDeadline的一半,这里设置为300秒(适配5分钟业务处理) private static final long ACK_EXTENSION_INTERVAL = 300; @ServiceActivator(inputChannel = "pubsubInputChannel") public void createVMListener(@Payload MyObjectThatNeedBeUnique payload, @Header(GcpPubSubHeaders.ORIGINAL_MESSAGE) BasicAcknowledgeablePubsubMessage message) throws IOException, ExecutionException, InterruptedException, TimeoutException { String messageId = message.getPubsubMessage().getMessageId(); log.info("Message arrived! Payload: {} | MessageId: {}", payload.toString(), messageId ); // 启动续期线程 ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); ScheduledFuture<?> future = scheduler.scheduleAtFixedRate(() -> { try { pubSubTemplate.modifyAckDeadline(subscriptionName, message.getAckId(), ACK_EXTENSION_INTERVAL); log.debug("Extended ack deadline for message: {}", messageId); } catch (Exception e) { log.error("Failed to extend ack deadline for message {}", messageId, e); } }, 0, ACK_EXTENSION_INTERVAL, TimeUnit.SECONDS); try { // 耗时5分钟的业务处理逻辑 // ...业务代码... // 处理完成后手动确认消息 message.ack(); log.info("Message {} processed successfully and acked", messageId); } catch (Exception e) { log.error("Failed to process message {}", messageId, e); // 处理失败时nack,让PubSub重新投递或进入死信队列 message.nack(); } finally { // 取消续期任务并关闭线程池 future.cancel(true); scheduler.shutdown(); } } }
2. 按需配置通道线程模型(可选)
如果业务不需要高并发处理,可将DirectChannel改为单线程执行,避免同一条消息被多线程同时处理:
@Bean public MessageChannel pubsubInputChannel() { DirectChannel channel = new DirectChannel(); channel.setTaskExecutor(Executors.newSingleThreadExecutor()); // 单线程执行器 return channel; }
注意:此方式会降低并发性能,仅适合业务本身无法并行的场景。
3. 实现业务层幂等性(最终兜底方案)
无论投递机制如何优化,网络波动、服务重启等极端情况仍可能导致重复投递,因此必须在业务层实现幂等性:
- 以消息ID或业务唯一标识(如userId)为键,处理前检查是否已处理过该消息
- 可使用数据库唯一约束、Redis缓存记录已处理消息ID等方式
示例:
// 假设使用Redis记录已处理消息 private final StringRedisTemplate redisTemplate; private static final String PROCESSED_PREFIX = "processed:messages:"; // 业务处理前检查 String processedKey = PROCESSED_PREFIX + messageId; if (redisTemplate.opsForValue().setIfAbsent(processedKey, "1", Duration.ofHours(24))) { // 未处理过,执行业务逻辑 } else { // 已处理过,直接确认消息并跳过 message.ack(); log.info("Message {} already processed, skipping", messageId); return; }
验证调整
修改后需确保:
- 业务处理过程中持续续期
ackDeadline,避免PubSub超时重投 - 仅在业务处理完成后才确认消息
- 业务层实现幂等性,兜底处理极端情况的重复投递
这样既能保留PubSubInboundChannelAdapter的流式拉取高效性,又能有效避免重复消费,满足恰好一次投递的要求。
内容的提问来源于stack exchange,提问作者Gustavo Steinmetz

