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

如何将基于Spring 5的ActiveMQ Classic迁移至Spring 6?

Spring 6 迁移 ActiveMQ Classic 请求/响应模式的解决方案

针对Spring 6移除org.springframework.jms.remoting.*包后的迁移需求,这里给你一套基于标准JMS API的请求/响应实现方案,完全适配Spring 6生态:

1. 依赖调整

先把Spring 5的JMS相关依赖换成Spring 6的spring-jms,同时保留ActiveMQ Classic的客户端依赖,Maven示例:

<dependency>
    <groupId>org.springframework</groupId>
    <artifactId>spring-jms</artifactId>
    <version>6.1.x</version> <!-- 用Spring 6最新稳定版 -->
</dependency>
<dependency>
    <groupId>org.apache.activemq</groupId>
    <artifactId>activemq-client</artifactId>
    <version>5.18.x</version> <!-- 匹配你的ActiveMQ Classic版本,建议升级到最新稳定版 -->
</dependency>

2. Web组件(请求发送端)改造

原来的JmsInvokerProxyFactoryBean这类远程调用组件要替换成手动构建请求消息,通过JMSReplyTo和JMSCorrelationID来关联请求与响应:

@Service
public class JmsRequestSender {
    @Autowired
    private JmsTemplate jmsTemplate;

    public String sendRequestAndWait(String requestQueue, String requestData, long timeout) throws JMSException {
        // 创建临时队列用于接收本次请求的响应
        Queue replyQueue = createTemporaryQueue();
        
        // 构建请求消息,指定响应队列和关联ID
        TextMessage requestMessage = createTextMessage(requestData);
        requestMessage.setJMSReplyTo(replyQueue);
        String correlationId = UUID.randomUUID().toString();
        requestMessage.setJMSCorrelationID(correlationId);

        // 发送请求到业务组件监听的队列
        jmsTemplate.send(requestQueue, session -> requestMessage);

        // 只接收和当前请求关联的响应,超时返回异常
        Message responseMessage = jmsTemplate.receiveSelected(replyQueue, 
            "JMSCorrelationID = '" + correlationId + "'", timeout);
        if (responseMessage instanceof TextMessage textMessage) {
            return textMessage.getText();
        }
        throw new RuntimeException("请求超时或未收到有效响应");
    }

    private Queue createTemporaryQueue() throws JMSException {
        try (Connection conn = jmsTemplate.getConnectionFactory().createConnection();
             Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE)) {
            return session.createTemporaryQueue();
        }
    }

    private TextMessage createTextMessage(String content) throws JMSException {
        try (Connection conn = jmsTemplate.getConnectionFactory().createConnection();
             Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE)) {
            return session.createTextMessage(content);
        }
    }
}

3. 业务组件(响应处理端)改造

替换原来的JmsInvokerServiceExporter,用标准的MessageListener处理请求,完成业务逻辑后发送响应到请求指定的队列:

@Component
public class JmsRequestListener implements MessageListener {
    @Autowired
    private JmsTemplate jmsTemplate;

    @Override
    public void onMessage(Message message) {
        try {
            if (!(message instanceof TextMessage requestMessage)) {
                throw new IllegalArgumentException("只支持TextMessage类型的请求");
            }

            // 执行业务逻辑处理请求
            String requestData = requestMessage.getText();
            String responseData = processRequest(requestData);

            // 获取请求指定的响应队列和关联ID
            Destination replyTo = requestMessage.getJMSReplyTo();
            String correlationId = requestMessage.getJMSCorrelationID();

            // 构建响应消息,关联原请求的ID
            TextMessage responseMessage = createTextMessage(responseData);
            responseMessage.setJMSCorrelationID(correlationId);

            // 发送响应到指定队列
            jmsTemplate.send(replyTo, session -> responseMessage);
        } catch (JMSException e) {
            throw new RuntimeException("处理JMS请求失败", e);
        }
    }

    private String processRequest(String requestData) {
        // 这里替换成你的实际业务逻辑
        return String.format("已处理请求:%s", requestData);
    }

    private TextMessage createTextMessage(String content) throws JMSException {
        try (Connection conn = jmsTemplate.getConnectionFactory().createConnection();
             Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE)) {
            return session.createTextMessage(content);
        }
    }
}

然后配置监听容器,让业务组件能监听到请求队列:

@Configuration
public class JmsConfig {
    @Value("${activemq.broker-url}")
    private String brokerUrl;

    @Bean
    public ConnectionFactory connectionFactory() {
        ActiveMQConnectionFactory connFactory = new ActiveMQConnectionFactory();
        connFactory.setBrokerURL(brokerUrl);
        // 按需配置用户名、密码、连接池等
        return connFactory;
    }

    @Bean
    public JmsTemplate jmsTemplate(ConnectionFactory connectionFactory) {
        JmsTemplate jmsTemplate = new JmsTemplate(connectionFactory);
        jmsTemplate.setSessionTransacted(true); // 按需开启事务
        return jmsTemplate;
    }

    @Bean
    public DefaultMessageListenerContainer requestListenerContainer(
            ConnectionFactory connectionFactory, JmsRequestListener listener) {
        DefaultMessageListenerContainer container = new DefaultMessageListenerContainer();
        container.setConnectionFactory(connectionFactory);
        container.setDestinationName("app-request-queue"); // 替换成你的请求队列名
        container.setMessageListener(listener);
        container.setSessionTransacted(true); // 和发送端保持事务一致性
        return container;
    }
}

4. 关键注意事项

  • 关联ID必须唯一:JMSCorrelationID是请求和响应的唯一标识,必须确保每次请求生成唯一值,避免响应串错
  • 临时队列vs固定队列:临时队列适合低并发场景,每次请求创建独立队列;高并发场景建议用固定响应队列,通过关联ID过滤响应
  • 超时机制:发送请求时一定要设置超时时间,避免线程无限阻塞
  • 事务配置:根据业务可靠性要求配置JMS事务,确保消息不丢失、不重复

可选替代方案:Spring Cloud Stream

如果你的系统已经接入Spring Cloud生态,也可以用Spring Cloud Stream的ActiveMQ binder来实现请求/响应,它封装了底层JMS细节,代码更简洁:

// 请求发送端
@Autowired
private StreamBridge streamBridge;

public String sendRequest(String requestData) {
    Message<String> requestMsg = MessageBuilder.withPayload(requestData)
            .setHeader(MessageHeaders.REPLY_CHANNEL, "app-response-topic")
            .build();
    return streamBridge.send("app-request-topic", requestMsg, String.class);
}

// 响应处理端
@Bean
public Consumer<Message<String>> requestProcessor() {
    return msg -> {
        String responseData = processRequest(msg.getPayload());
        streamBridge.send(msg.getHeaders().get(MessageHeaders.REPLY_CHANNEL, String.class), responseData);
    };
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 11:27:47