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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 13:24:56