Spring Boot集成Jakarta JMS:如何实现消息即时发送?
问题:JMS消息无法实时发送,需等待事务结束才批量投递
我需要在不等待事务结束的情况下发送消息,但尝试过template.setSessionTransacted(false);及设置不同确认模式等操作后均未生效。
使用的JMS依赖
implementation 'org.apache.activemq:artemis-jakarta-client:2.38.0' implementation 'org.apache.activemq:artemis-core-client:2.38.0' implementation 'org.apache.activemq:artemis-commons:2.38.0' implementation 'org.springframework:spring-jms:6.1.14'
JMS配置代码
@Bean public ConnectionFactory connectionFactory() { ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(); connectionFactory.setBrokerURL(format("tcp://%s:%s", host, port)); connectionFactory.setUser(user); connectionFactory.setPassword(pass); return connectionFactory; } @Bean public JmsListenerContainerFactory<?> myFactory(ConnectionFactory connectionFactory, DefaultJmsListenerContainerFactoryConfigurer configurer) { Logger logger = (Logger) LoggerFactory.getLogger("org.springframework.jms"); logger.setLevel(Level.ERROR); DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); factory.setErrorHandler(new CustomErrorHandler()); configurer.configure(factory, connectionFactory); return factory; } @Bean public JmsTemplate jmsAnycastTemplate(ConnectionFactory connectionFactory) { return new JmsTemplate(connectionFactory); }
Worker客户端代码
public static void doAllWork(int moduleId, MqProducer mqProducer, MqMessage message){ mqProducer.sendAcceptToWork(moduleId, message); Work.doSomeWork(moduleId, mqProducer); mqProducer.sendEndWork(moduleId, message, 200L); } @SneakyThrows public static void doSomeWork(int moduleId, MqProducer mqProducer){ int secToWork = ThreadLocalRandom.current().nextInt(1, 10); LocalDateTime workstarted = LocalDateTime.now(); while (Duration.between(workstarted, LocalDateTime.now()).getSeconds() < secToWork) { Thread.sleep(1000); System.out.printf("%s [api] WORK_ENDED message received%n", LocalDateTime.now()); mqProducer.sendMessage("api", new MqMessage(IN_PROGRESS, moduleId, 200)); } }
Worker客户端输出
sendAcceptToWork 3 0 2024-12-19T13:03:45.046627435 Working... 2024-12-19T13:03:46.057809869 Working... 2024-12-19T13:03:47.065647593 Working... 2024-12-19T13:03:48.074799444 Working... 2024-12-19T13:03:49.080706764 Working... 2024-12-19T13:03:50.088089181 Working... 2024-12-19T13:03:51.095799844 Working... 2024-12-19T13:03:52.104492990 Working... 2024-12-19T13:03:53.113162971 Working... Work done
API客户端代码
@SneakyThrows @Override public void run(MqMessage message) { if (message.getMqOperationCode() == MqOperationCode.WORK_STARTED){ System.out.println(); System.out.printf("%s [api] WORK_STARTED message received%n", LocalDateTime.now()); return; } if (message.getMqOperationCode() == MqOperationCode.IN_PROGRESS){ System.out.printf("%s [api] IN_PROGRESS message received%n", LocalDateTime.now()); return; } if (message.getMqOperationCode() == MqOperationCode.WORK_ENDED){ System.out.printf("%s [api] WORK_ENDED message received%n", LocalDateTime.now()); return; } System.out.printf("%s [api] message received : \n%s%n", LocalDateTime.now(), message); }
API客户端输出
2024-12-19T13:03:53.141036336 [api] WORK_STARTED message received 2024-12-19T13:03:53.146047337 [api] IN_PROGRESS message received 2024-12-19T13:03:53.150053390 [api] IN_PROGRESS message received 2024-12-19T13:03:53.153668788 [api] IN_PROGRESS message received 2024-12-19T13:03:53.157834576 [api] IN_PROGRESS message received 2024-12-19T13:03:53.161060836 [api] IN_PROGRESS message received 2024-12-19T13:03:53.163221293 [api] IN_PROGRESS message received 2024-12-19T13:03:53.165163888 [api] IN_PROGRESS message received 2024-12-19T13:03:53.166759323 [api] IN_PROGRESS message received 2024-12-19T13:03:53.168365120 [api] IN_PROGRESS message received 2024-12-19T13:03:53.169784977 [api] WORK_ENDED message received
如上述输出所示,Worker客户端的所有消息会一次性发送,但我需要实现每条消息无需等待事务结束即可实时投递。
解决方案
1. 排查事务上下文
首先确认doAllWork或其调用链上是否存在@Transactional注解,Spring事务会自动将JMS操作绑定到当前事务,导致消息延迟到事务提交时批量发送。
2. 配置JmsTemplate脱离事务
修改JmsTemplate配置,强制其不参与当前事务:
@Bean public JmsTemplate jmsAnycastTemplate(ConnectionFactory connectionFactory) { JmsTemplate template = new JmsTemplate(connectionFactory); template.setSessionTransacted(false); template.setSessionAcknowledgeMode(Session.AUTO_ACKNOWLEDGE); // 关键:移除事务管理器绑定,避免参与当前事务 template.setTransactionManager(null); return template; }
3. 使用独立非事务JmsTemplate
如果上述配置无效,创建完全独立的非事务JmsTemplate:
@Bean public JmsTemplate nonTransactionalJmsTemplate() { ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(); connectionFactory.setBrokerURL(format("tcp://%s:%s", host, port)); connectionFactory.setUser(user); connectionFactory.setPassword(pass); JmsTemplate template = new JmsTemplate(connectionFactory); template.setSessionTransacted(false); template.setSessionAcknowledgeMode(Session.AUTO_ACKNOWLEDGE); return template; }
在MqProducer中注入该模板发送实时消息。
4. 手动控制JMS会话发送
直接使用原生JMS API,手动创建非事务会话确保消息立即发送:
@SneakyThrows public void sendMessage(String destination, MqMessage message) { try (Connection conn = connectionFactory.createConnection(); Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE); MessageProducer producer = session.createProducer(session.createQueue(destination))) { conn.start(); ObjectMessage objMsg = session.createObjectMessage(message); producer.send(objMsg); // AUTO_ACKNOWLEDGE模式下发送即完成确认,无需额外提交 } }
5. 调整监听器容器事务配置
如果消息触发自JMS监听器,检查监听器容器的事务设置:
@Bean public JmsListenerContainerFactory<?> myFactory(ConnectionFactory connectionFactory, DefaultJmsListenerContainerFactoryConfigurer configurer) { DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); factory.setErrorHandler(new CustomErrorHandler()); // 禁用监听器事务,避免影响发送操作 factory.setSessionTransacted(false); factory.setSessionAcknowledgeMode(Session.AUTO_ACKNOWLEDGE); configurer.configure(factory, connectionFactory); return factory; }
内容的提问来源于stack exchange,提问作者mamol
相关产品推荐
相关产品推荐

