Azure Service Bus JMS Spring长任务超时重复消费问题咨询
问题描述
Azure Service Bus会在Azure Blob上传新XML文件时接收消息,Spring微服务通过@JmsListener监听队列处理消息。小文件处理正常,但30GB+大文件处理耗时超10分钟,监听器5分钟超时后消息重回队列,导致重复消费。因业务限制需串行处理文件,咨询以下问题:
- 是否可延长超时时间,确保监听器完成首个任务?
- 若将长任务放入异步方法,如何避免监听器接收新消息?
- 若监听器已接收新消息,如何优雅地确认或放回队列?
当前实现代码
@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
相关产品推荐
相关产品推荐

