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

Spring Boot中PubSubInboundChannelAdapter消息重复问题咨询

Spring Boot + Spring Integration PubSub 恰好一次投递问题排查与解决

问题背景

我正在开发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个线程同时处理。

疑问

  1. 为何出现消息重复?
  2. 如何在保留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;
    }

解答

问题原因

  1. AUTO_ACK模式时机错误:PubSubInboundChannelAdapter采用异步拉取,AUTO_ACK模式下会在消息发送到通道后立即确认,而非等待业务处理完成。如果业务处理耗时超过PubSub的ackDeadline(即使设置了max-ack-extension-period,客户端无主动续期时,PubSub仍会判定消息处理超时),就会触发重新投递。
  2. 流式拉取的并发特性:PubSubInboundChannelAdapter基于流式拉取,默认一次性拉取多条消息分发到通道,而DirectChannel默认使用多线程线程池处理消息。若消息提前确认但业务处理超时/失败,PubSub重投后会被线程池的多个线程同时处理,导致重复消费。
  3. 缺少主动续期机制: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 11:25:30