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

Azure Service Bus JMS Spring长任务超时重复消费问题咨询

问题描述

Azure Service Bus会在Azure Blob上传新XML文件时接收消息,Spring微服务通过@JmsListener监听队列处理消息。小文件处理正常,但30GB+大文件处理耗时超10分钟,监听器5分钟超时后消息重回队列,导致重复消费。因业务限制需串行处理文件,咨询以下问题:

  1. 是否可延长超时时间,确保监听器完成首个任务?
  2. 若将长任务放入异步方法,如何避免监听器接收新消息?
  3. 若监听器已接收新消息,如何优雅地确认或放回队列?

当前实现代码

@JmsListener(
    destination = "file-name-queue",
    containerFactory = "jmsListenerContainerFactory"
)
public void processMessages(JmsMessage message) throws JMSException{
    
    String fileName = message.getBody(String.class);
    log.info("Process The file {}",fileName);
    // 这是耗时很长的处理逻辑
    processFullFile(fileName);

    log.info("Process finished for The file {}",fileName);

}

计划修改的代码

@JmsListener(
    destination = "file-name-queue",
    containerFactory = "jmsListenerContainerFactory"
)
public void processMessages(JmsMessage message) throws JMSException{

    String fileName = message.getBody(String.class);
    log.info("Received a new file {}",fileName);
    // 检查是否有文件正在处理
    if(fileInProgress()) {
        log.info("Existing file already in progress");
        
        //TODO: 如何让消息重回队列?
        // 尝试过等待一段时间后抛出异常,虽然能让消息回到队列,但会排在队尾,且不够优雅
        
    } else {
        log.info("Process The file {}",fileName);
        // 这是耗时很长的处理逻辑
        aysncProcessFullFile(fileName);
    }
}

当前Spring配置

spring: 
  jms: 
    servicebus: 
      enabled: true 
      connection-string: Endpoint=sb://XXXXXX.servicebus.windows.net/;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=SHARED_ACCESS_KEY 
      idle-timeout: 2000000 
      pricing-tier: Standard

异常信息

10:37:56.246 [org.springframework.jms.JmsListenerEndpointContainer#0-1] INFO  c.h.s.c.config.MessageQueueConfig - Under process
10:37:56.286 [org.springframework.jms.JmsListenerEndpointContainer#0-1] INFO  o.a.q.jms.JmsLocalTransactionContext - Commit failed for transaction: TX:ID:2897b0c2-8f4e-41d3-9122-f1f89c9f8b83:1:115
10:37:56.287 [org.springframework.jms.JmsListenerEndpointContainer#0-1] INFO  o.s.j.c.CachingConnectionFactory - Encountered a JMSException - resetting the underlying JMS Connection
jakarta.jms.TransactionRolledBackException: Transaction is not declared Reference:1cbe9d94-65c4-4a0b-a58d-29bbdf2d616d, TrackingId:6b39f9d4-c914-4ec8-b764-53cbc484b136_G1, SystemTracker:gtm, Timestamp:2023-08-28T08:37:56 [condition = amqp:transaction:unknown-id]
解决方案

1. 延长超时时间确保监听器完成任务

可以通过调整Azure Service Bus锁配置和Spring JMS参数实现:

  • Azure Service Bus端:将队列的Lock Duration设置为标准层允许的最大值(5分钟)。
  • Spring配置端:开启锁自动续期,并调整事务超时时间覆盖文件处理时长:
    spring:
      jms:
        servicebus:
          enabled: true
          connection-string: Endpoint=sb://XXXXXX.servicebus.windows.net/;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=SHARED_ACCESS_KEY
          idle-timeout: 2000000
          pricing-tier: Standard
          lock-renewal-duration: 3600000  # 1小时,确保覆盖文件处理时长
        listener:
          acknowledge-mode: AUTO
          transaction:
            timeout: 3600  # 事务超时时间(秒)
    
    注意:锁续期会自动维持消息的锁定状态,直到任务完成或续期超时,避免消息因锁过期重回队列。

2. 异步处理时避免监听器接收新消息

结合串行需求,推荐两种方案:

  • 方案一:暂停/恢复监听器容器
    给@JmsListener添加id属性,通过JmsListenerEndpointRegistry控制容器启停:

    @Autowired
    private JmsListenerEndpointRegistry registry;
    
    @JmsListener(
        id = "file-name-queue",
        destination = "file-name-queue",
        containerFactory = "jmsListenerContainerFactory"
    )
    public void processMessages(JmsMessage message) throws JMSException{
        String fileName = message.getBody(String.class);
        log.info("Received a new file {}",fileName);
        
        if(!fileInProgress()) {
            log.info("Process The file {}",fileName);
            asyncProcessFullFile(fileName);
        }
    }
    
    private void asyncProcessFullFile(String fileName) {
        // 暂停监听器,不再接收新消息
        registry.getListenerContainer("file-name-queue").stop();
        try {
            processFullFile(fileName);
        } finally {
            // 任务完成后恢复监听器
            registry.getListenerContainer("file-name-queue").start();
        }
    }
    
  • 方案二:手动确认+单并发
    将监听器设置为MANUAL确认模式,并发数保持1,异步任务完成后再手动确认消息:

    @JmsListener(
        id = "file-name-queue",
        destination = "file-name-queue",
        containerFactory = "jmsListenerContainerFactory",
        acknowledgeMode = "MANUAL"
    )
    public void processMessages(JmsMessage message) throws JMSException{
        String fileName = message.getBody(String.class);
        log.info("Received a new file {}",fileName);
    
        // 异步执行长任务,完成后确认消息
        CompletableFuture.runAsync(() -> {
            try {
                processFullFile(fileName);
                message.acknowledge(); // 任务完成确认消息
            } catch (Exception e) {
                message.setJMSRedelivered(true);
                // 处理失败,让消息重回队列
            }
        });
    }
    

    此方式下,监听器会持有当前消息锁直到手动确认,因并发数为1,不会接收新消息。

3. 优雅将已接收消息放回队列

推荐两种优雅处理方式:

  • 手动会话恢复
    在MANUAL确认模式下,通过Session将消息立即放回队列:

    @JmsListener(
        id = "file-name-queue",
        destination = "file-name-queue",
        containerFactory = "jmsListenerContainerFactory",
        acknowledgeMode = "MANUAL"
    )
    public void processMessages(JmsMessage message, Session session) throws JMSException{
        String fileName = message.getBody(String.class);
        log.info("Received a new file {}",fileName);
    
        if(fileInProgress()) {
            log.info("Existing file already in progress, putting message back");
            session.recover(); // 消息立即重回队列
            return;
        }
        // 处理任务...
    }
    
  • 延迟重试调度
    利用Azure Service Bus的延迟入队特性,将消息安排到未来时间再处理:

    @JmsListener(
        id = "file-name-queue",
        destination = "file-name-queue",
        containerFactory = "jmsListenerContainerFactory",
        acknowledgeMode = "MANUAL"
    )
    public void processMessages(JmsMessage message, Session session) throws JMSException{
        String fileName = message.getBody(String.class);
        log.info("Received a new file {}",fileName);
    
        if(fileInProgress()) {
            log.info("Existing file in progress, scheduling message for later");
            // 延迟1小时重试
            long delayMillis = 3600000;
            message.setLongProperty("ScheduledEnqueueTimeUtc", System.currentTimeMillis() + delayMillis);
            
            message.acknowledge(); // 确认当前消息
            // 重新发送到原队列
            session.createProducer(session.createQueue("file-name-queue")).send(message);
            return;
        }
        // 处理任务...
    }
    

    此方式避免短时间内重复消费同一消息,更贴合串行处理的业务节奏。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 00:39:53