如何将基于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
相关产品推荐
相关产品推荐

